1use base64::{Engine as _, engine::general_purpose};
2use chrono::Utc;
3use futures::{StreamExt as _, stream::FuturesOrdered};
4use haste_fhir_client::FHIRClient;
5use haste_fhir_generated_ops::generated::ViewDefinitionRun;
6use haste_fhir_model::r4::{
7 self,
8 generated::{
9 resources::{Binary, Bundle, Resource, ResourceType, ViewDefinition, ViewDefinitionSelect},
10 terminology::{BoundCode, IssueType, OutputFormatCodes},
11 types::{FHIRBase64Binary, FHIRBoolean, Reference},
12 },
13};
14use haste_fhir_operation_error::OperationOutcomeError;
15use haste_fhirpath::{Config, FPEngine};
16use haste_reflect::MetaValue;
17use itertools::Itertools as _;
18use ordermap::OrderMap;
19use serde::{Deserialize, Serialize};
20use std::{borrow::Cow, collections::HashMap, sync::Arc};
21
22use crate::conversions::primitives::PrimitiveValue;
23
24mod compartment;
25mod conversions;
26mod output;
27
28fn reference_value(reference: &Reference) -> Result<String, OperationOutcomeError> {
29 reference
30 .reference
31 .as_ref()
32 .and_then(|r| r.value.clone())
33 .ok_or_else(|| {
34 OperationOutcomeError::error(
35 IssueType::invalid(),
36 "Reference.reference is required".to_string(),
37 )
38 })
39}
40
41async fn resolve_patient_references<
47 CTX: Send + Sync + Clone + 'static,
48 Client: FHIRClient<CTX, OperationOutcomeError> + Send + Sync + 'static,
49>(
50 context: CTX,
51 client: &Client,
52 input: &ViewDefinitionRun::Input,
53) -> Result<Vec<String>, OperationOutcomeError> {
54 let mut references = Vec::new();
55
56 if let Some(patient) = input.patient.as_ref() {
57 references.push(reference_value(patient)?);
58 }
59
60 for group_reference in input.group.as_ref().into_iter().flatten() {
61 let group_reference_value = reference_value(group_reference)?;
62 let group_id = group_reference_value
63 .rsplit('/')
64 .next()
65 .unwrap_or(&group_reference_value)
66 .to_string();
67
68 let group = client
69 .read(context.clone(), ResourceType::Group, group_id.clone())
70 .await?
71 .ok_or_else(|| {
72 OperationOutcomeError::error(
73 IssueType::not_found(),
74 format!("Group not found with id '{group_id}'"),
75 )
76 })?;
77
78 let Resource::Group(group) = group else {
79 return Err(OperationOutcomeError::error(
80 IssueType::invalid(),
81 format!("Reference '{group_reference_value}' does not point to a Group resource"),
82 ));
83 };
84
85 for member in group.member.into_iter().flatten() {
86 let entity_reference = reference_value(&member.entity)?;
87 if entity_reference.starts_with("Patient/") {
88 references.push(entity_reference);
89 }
90 }
91 }
92
93 references.sort_unstable();
94 references.dedup();
95
96 Ok(references)
97}
98
99async fn resolve_view_definition<
100 'a,
101 CTX: Send + Sync + Clone + 'static,
102 Client: FHIRClient<CTX, OperationOutcomeError> + Send + Sync + 'static,
103>(
104 context: CTX,
105 client: &Client,
106 input: &'a ViewDefinitionRun::Input,
107) -> Result<Cow<'a, ViewDefinition>, OperationOutcomeError> {
108 if let Some(view_definition) = &input.viewResource {
109 Ok(Cow::Borrowed(view_definition))
110 } else if let Some(view_definition_reference) = input.viewReference.as_ref() {
111 let view_definition_reference = view_definition_reference
112 .reference
113 .as_ref()
114 .ok_or_else(|| {
115 OperationOutcomeError::error(
116 IssueType::invalid(),
117 "viewReference.reference is required".to_string(),
118 )
119 })?
120 .value
121 .as_ref()
122 .ok_or_else(|| {
123 OperationOutcomeError::error(
124 IssueType::invalid(),
125 "viewReference.reference.value is required".to_string(),
126 )
127 })?;
128
129 let reference_pieces = view_definition_reference.split('/').collect::<Vec<_>>();
130
131 let view_definition_id = reference_pieces
132 .last()
133 .ok_or_else(|| {
134 OperationOutcomeError::error(
135 IssueType::invalid(),
136 "Invalid viewReference.reference format".to_string(),
137 )
138 })?
139 .to_string();
140
141 let result = client
142 .read(
143 context,
144 ResourceType::ViewDefinition,
145 view_definition_id.clone(),
146 )
147 .await?
148 .ok_or_else(|| {
149 OperationOutcomeError::error(
150 IssueType::not_found(),
151 format!("ViewDefinition not found with id '{view_definition_id:?}'"),
152 )
153 })?;
154
155 if let Resource::ViewDefinition(view_definition) = result {
156 Ok(Cow::Owned(view_definition))
157 } else {
158 Err(OperationOutcomeError::error(
159 IssueType::invalid(),
160 "Referenced resource is not a ViewDefinition".to_string(),
161 ))
162 }
163 } else {
164 Err(OperationOutcomeError::error(
165 IssueType::invalid(),
166 "Either viewResource or viewReference must be provided".to_string(),
167 ))
168 }
169}
170
171const HISTORY_PAGE_SIZE: u32 = 1000;
174
175const MAX_RESOURCES_PER_RUN: usize = 50_000;
183
184fn requested_limit(input: &ViewDefinitionRun::Input) -> Option<usize> {
188 input
189 .limit
190 .as_ref()
191 .and_then(|limit| limit.value)
192 .and_then(|limit| usize::try_from(limit).ok())
193}
194
195async fn get_resources_to_process<
196 CTX: Send + Sync + Clone + 'static,
197 Client: FHIRClient<CTX, OperationOutcomeError> + Send + Sync + 'static,
198>(
199 context: CTX,
200 client: &Client,
201 view_definition: &ViewDefinition,
202 input: &ViewDefinitionRun::Input,
203) -> Result<Vec<Resource>, OperationOutcomeError> {
204 if let Some(input_resources) = input.resource.clone() {
205 return Ok(input_resources);
206 }
207
208 let Some(resource_type) = view_definition.resource.as_str() else {
209 return Err(OperationOutcomeError::error(
210 IssueType::invalid(),
211 "ViewDefinition.resource is required".to_string(),
212 ));
213 };
214
215 let resource_type = ResourceType::try_from(resource_type).map_err(|e| {
216 OperationOutcomeError::error(IssueType::invalid(), format!("Invalid resource type: {e}"))
217 })?;
218
219 let since_instant = input.since.as_ref().and_then(|since| since.value.clone());
220
221 let patient_references = resolve_patient_references(context.clone(), client, input).await?;
222
223 if !patient_references.is_empty() {
224 let last_updated_filter = since_instant.as_ref().map(|since| format!("gt{since:?}"));
225
226 return compartment::resources_for_patients(
227 context,
228 client,
229 resource_type,
230 &patient_references,
231 last_updated_filter.as_deref(),
232 )
233 .await;
234 }
235
236 let since = since_instant.unwrap_or(r4::datetime::Instant::Iso8601(Utc::now()));
237
238 let resource_limit = requested_limit(input).map_or(MAX_RESOURCES_PER_RUN, |limit| {
239 limit.min(MAX_RESOURCES_PER_RUN)
240 });
241
242 let mut combined: Option<Bundle> = None;
243 let mut offset = 0u32;
244 let mut total_fetched = 0usize;
245
246 loop {
247 let mut page = client
248 .history_type(
249 context.clone(),
250 resource_type.clone(),
251 vec![
252 ("_since".to_string(), vec![since.to_string()]),
253 ("_count".to_string(), vec![HISTORY_PAGE_SIZE.to_string()]),
254 ("_offset".to_string(), vec![offset.to_string()]),
255 ]
256 .into(),
257 )
258 .await?;
259
260 let entry_count = page.entry.as_ref().map_or(0, Vec::len);
261 total_fetched += entry_count;
262
263 match combined.as_mut() {
264 Some(accumulated) => accumulated
265 .entry
266 .get_or_insert_with(Vec::new)
267 .extend(page.entry.take().into_iter().flatten()),
268 None => combined = Some(page),
269 }
270
271 if entry_count < HISTORY_PAGE_SIZE as usize || total_fetched >= resource_limit {
272 break;
273 }
274
275 offset += HISTORY_PAGE_SIZE;
276 }
277
278 let mut combined = combined.unwrap_or_default();
279
280 if let Some(entries) = combined.entry.as_mut() {
283 entries.truncate(resource_limit);
284 }
285
286 Ok(vec![Resource::Bundle(combined)])
287}
288
289fn build_hashmap_fp_variables(viewdefinition: &ViewDefinition) -> HashMap<String, &dyn MetaValue> {
290 let mut hashmap = HashMap::new();
291
292 if let Some(constants) = &viewdefinition.constant {
293 for constant in constants {
294 if let Some(name) = &constant.name.value.as_ref() {
295 hashmap.insert((*name).clone(), &constant.value as &dyn MetaValue);
296 }
297 }
298 }
299
300 hashmap
301}
302
303fn cartesian_product(
304 select_statement_results: Vec<Vec<OrderMap<String, OutputResults>>>,
305) -> Vec<OrderMap<String, OutputResults>> {
306 let mut output_results = Vec::new();
307
308 for combination in select_statement_results
309 .into_iter()
310 .multi_cartesian_product()
311 {
312 let mut combined_result = OrderMap::new();
313
314 for result in combination {
315 for (key, value) in result {
316 combined_result.insert(key, value);
317 }
318 }
319
320 output_results.push(combined_result);
321 }
322
323 output_results
324}
325
326#[derive(Debug, Clone, Deserialize, Serialize)]
328#[serde(untagged)]
329enum OutputResults {
330 Scalar(Option<PrimitiveValue>),
331 Collection(Vec<Option<PrimitiveValue>>),
332}
333
334async fn process_resource<
335 CTX: Send + Sync + Clone + 'static,
336 Client: FHIRClient<CTX, OperationOutcomeError> + Send + Sync + 'static,
337>(
338 _context: CTX,
339 _client: Arc<Client>,
340 variables: Arc<HashMap<String, &dyn MetaValue>>,
341 view_definition: &ViewDefinition,
342 input: Resource,
343) -> Result<Vec<OrderMap<String, OutputResults>>, OperationOutcomeError> {
344 let fp_engine = FPEngine::new();
345
346 let mut select_statement_results = Vec::with_capacity(view_definition.select.len());
347
348 for select_statement in &view_definition.select {
349 let fp_config = Arc::new(
350 Config::builder()
351 .with_variable_resolver(haste_fhirpath::ExternalConstantResolver::Variable(
352 variables.clone(),
353 ))
354 .with_resource_id(input.id().clone().unwrap_or_default()),
355 );
356
357 let (iterable_context, set_null) =
358 build_iterable_context(&fp_engine, fp_config.clone(), select_statement, &input).await?;
359
360 let select_results = process_select_statement(
361 &fp_engine,
362 fp_config,
363 select_statement,
364 &input,
365 iterable_context,
366 set_null,
367 )
368 .await?;
369
370 select_statement_results.push(select_results);
371 }
372
373 let output_results = cartesian_product(select_statement_results);
374
375 Ok(output_results)
376}
377
378async fn build_iterable_context<'a>(
379 fp_engine: &FPEngine,
380 fp_config: Arc<Config<'a>>,
381 select_statement: &'a ViewDefinitionSelect,
382 input: &'a Resource,
383) -> Result<(Option<Vec<haste_fhirpath::Context<'a>>>, bool), OperationOutcomeError> {
384 let mut iterable_context = None;
385 let mut set_null = false;
386
387 if let Some(for_each_fp) = select_statement
388 .forEach
389 .as_ref()
390 .and_then(|f| f.value.as_ref())
391 {
392 iterable_context = Some(vec![
393 fp_engine
394 .evaluate_with_config(for_each_fp, vec![input], fp_config.clone())
395 .await
396 .map_err(|e| {
397 OperationOutcomeError::error(
398 IssueType::exception(),
399 format!("Error evaluating forEach expression: {e}"),
400 )
401 })?,
402 ]);
403 } else if let Some(for_each_or_null_fp) = select_statement
404 .forEachOrNull
405 .as_ref()
406 .and_then(|f| f.value.as_ref())
407 {
408 iterable_context = Some(vec![
409 fp_engine
410 .evaluate_with_config(for_each_or_null_fp, vec![input], fp_config.clone())
411 .await
412 .map_err(|e| {
413 OperationOutcomeError::error(
414 IssueType::exception(),
415 format!("Error evaluating forEachOrNull expression: {e}"),
416 )
417 })?,
418 ]);
419
420 set_null = true;
421 } else if let Some(repeat) = select_statement
422 .repeat
423 .as_ref()
424 .map(|r| r.iter().filter_map(|r| r.value.as_ref()))
425 {
426 let mut repeat_fps = vec![];
427
428 for repeat_fp in repeat {
429 let repeat = format!("$this.repeat({repeat_fp})");
430
431 repeat_fps.push(
432 fp_engine
433 .evaluate_with_config(&repeat, vec![input], fp_config.clone())
434 .await
435 .map_err(|e| {
436 OperationOutcomeError::error(
437 IssueType::exception(),
438 format!("Error evaluating repeat expression: {e}"),
439 )
440 })?,
441 );
442 }
443
444 iterable_context = Some(repeat_fps);
445 }
446
447 Ok((iterable_context, set_null))
448}
449
450async fn process_select_statement<'a>(
451 fp_engine: &FPEngine,
452 fp_config: Arc<Config<'a>>,
453 select_statement: &'a ViewDefinitionSelect,
454 input: &'a Resource,
455 iterable_context: Option<Vec<haste_fhirpath::Context<'a>>>,
456 set_null: bool,
457) -> Result<Vec<OrderMap<String, OutputResults>>, OperationOutcomeError> {
458 let select_context: Vec<&dyn MetaValue> = if let Some(iterable) = iterable_context.as_ref() {
459 iterable
460 .iter()
461 .flat_map(haste_fhirpath::Context::iter)
462 .collect()
463 } else {
464 vec![input]
465 };
466
467 let mut select_results = Vec::with_capacity(select_context.len());
468
469 if set_null && select_context.is_empty() {
470 let output_result = build_null_result(select_statement)?;
471 select_results.push(output_result);
472 }
473
474 for context in select_context {
475 let output_result =
476 process_select_context(fp_engine, fp_config.clone(), select_statement, context).await?;
477
478 select_results.push(output_result);
479 }
480
481 Ok(select_results)
482}
483
484fn build_null_result(
485 select_statement: &ViewDefinitionSelect,
486) -> Result<OrderMap<String, OutputResults>, OperationOutcomeError> {
487 let mut output_result = OrderMap::new();
488
489 for column in select_statement.column.as_ref().into_iter().flatten() {
490 let Some(name) = column.name.value.as_deref() else {
491 return Err(OperationOutcomeError::error(
492 IssueType::invalid(),
493 "Column name is required".to_string(),
494 ));
495 };
496
497 output_result.insert(name.to_string(), OutputResults::Scalar(None));
498 }
499
500 Ok(output_result)
501}
502
503async fn process_select_context<'a>(
504 fp_engine: &FPEngine,
505 fp_config: Arc<Config<'a>>,
506 select_statement: &'a ViewDefinitionSelect,
507 context: &'a dyn MetaValue,
508) -> Result<OrderMap<String, OutputResults>, OperationOutcomeError> {
509 let mut output_result = OrderMap::new();
510
511 for column in select_statement.column.as_ref().into_iter().flatten() {
512 let Some(path) = column.path.value.as_deref() else {
513 return Err(OperationOutcomeError::error(
514 IssueType::invalid(),
515 "Column path is required".to_string(),
516 ));
517 };
518
519 let Some(name) = column.name.value.as_deref() else {
520 return Err(OperationOutcomeError::error(
521 IssueType::invalid(),
522 "Column name is required".to_string(),
523 ));
524 };
525
526 let result = fp_engine
527 .evaluate_with_config(path, vec![context], fp_config.clone())
528 .await
529 .map_err(|e| {
530 OperationOutcomeError::error(
531 IssueType::exception(),
532 format!("Error evaluating expression: {e}"),
533 )
534 })?;
535
536 let column_type = column
537 .type_
538 .as_ref()
539 .and_then(|t| t.value.as_deref())
540 .unwrap_or_else(|| {
541 result
542 .iter()
543 .next()
544 .map_or("string", haste_reflect::MetaValue::fhir_type)
545 });
546
547 let mut column_result = result
548 .iter()
549 .map(|value| conversions::primitives::convert_meta_value(column_type, value))
550 .collect::<Result<Vec<Option<PrimitiveValue>>, OperationOutcomeError>>()?;
551
552 let is_collection = column
553 .collection
554 .as_ref()
555 .and_then(|c| c.value)
556 .unwrap_or(false);
557
558 let insert_value = if is_collection {
559 OutputResults::Collection(column_result)
560 } else {
561 if column_result.len() > 1 {
562 return Err(OperationOutcomeError::error(
563 IssueType::invalid(),
564 "Column result is a collection but the column is not marked as a collection"
565 .to_string(),
566 ));
567 }
568
569 let mut singular_value = None;
570
571 if let Some(first_value) = column_result.get_mut(0) {
572 std::mem::swap(&mut singular_value, first_value);
573 }
574
575 OutputResults::Scalar(singular_value)
576 };
577
578 output_result.insert(name.to_string(), insert_value);
579 }
580
581 Ok(output_result)
582}
583
584fn flatten_results(resource: Vec<Resource>) -> Vec<Resource> {
585 let mut resources = Vec::new();
586 for resource in resource {
587 match resource {
588 Resource::Bundle(bundle) => {
589 for entry in bundle.entry.into_iter().flatten() {
590 if let Some(resource) = entry.resource {
591 resources.push(*resource);
592 }
593 }
594 }
595 _ => {
596 resources.push(resource);
597 }
598 }
599 }
600
601 resources
602}
603
604async fn passes_where_clauses(
605 fp_engine: &FPEngine,
606 variables: Arc<HashMap<String, &dyn MetaValue>>,
607 where_clauses: &[&str],
608 resource: &Resource,
609) -> Result<bool, OperationOutcomeError> {
610 for where_clause in where_clauses {
611 let result = fp_engine
612 .evaluate_with_config(
613 where_clause,
614 vec![resource],
615 Arc::new(Config::builder().with_variable_resolver(
616 haste_fhirpath::ExternalConstantResolver::Variable(variables.clone()),
617 )),
618 )
619 .await
620 .map_err(|e| {
621 OperationOutcomeError::error(
622 IssueType::exception(),
623 format!("Error evaluating where clause expression: {e}"),
624 )
625 })?;
626
627 let bool_results = result
628 .iter()
629 .map(|v| match v.fhir_type() {
630 "boolean" => Ok(v
631 .as_any()
632 .downcast_ref::<FHIRBoolean>()
633 .and_then(|b| b.value.as_ref())
634 .unwrap_or(&false)),
635 "http://hl7.org/fhirpath/System.Boolean" => {
636 Ok(v.as_any().downcast_ref::<bool>().unwrap_or(&false))
637 }
638 _ => Err(OperationOutcomeError::error(
639 IssueType::invalid(),
640 format!(
641 "Where clause expression must evaluate to a boolean, got: {}",
642 v.fhir_type()
643 ),
644 )),
645 })
646 .collect::<Result<Vec<_>, _>>()?;
647
648 if bool_results.iter().any(|v| !**v) {
649 return Ok(false);
650 }
651 }
652
653 Ok(true)
654}
655
656async fn process_view_definition<
657 CTX: Send + Sync + Clone + 'static,
658 Client: FHIRClient<CTX, OperationOutcomeError> + Send + Sync + 'static,
659>(
660 context: CTX,
661 output_format: &BoundCode<OutputFormatCodes>,
662 client: Arc<Client>,
663 view_definition: &ViewDefinition,
664 input: &ViewDefinitionRun::Input,
665) -> Result<Binary, OperationOutcomeError> {
666 let variables = Arc::new(build_hashmap_fp_variables(view_definition));
667 let limit = requested_limit(input);
668
669 let input_ = flatten_results(
670 get_resources_to_process(context.clone(), client.as_ref(), view_definition, input).await?,
671 );
672
673 let mut tasks = FuturesOrdered::new();
674
675 let where_clauses = view_definition
676 .where_
677 .as_ref()
678 .map_or_else(|| Cow::Owned(Vec::new()), Cow::Borrowed);
679
680 let where_fp_clauses = where_clauses
681 .iter()
682 .filter_map(|w| w.path.value.as_deref())
683 .collect::<Vec<_>>();
684
685 for resource in input_ {
686 if passes_where_clauses(
687 &FPEngine::new(),
688 variables.clone(),
689 where_fp_clauses.as_slice(),
690 &resource,
691 )
692 .await?
693 {
694 tasks.push_back(async {
695 process_resource(
696 context.clone(),
697 client.clone(),
698 variables.clone(),
699 view_definition,
700 resource,
701 )
702 .await
703 });
704 }
705 }
706
707 let mut results = Vec::with_capacity(tasks.len());
708
709 while let Some(result) = tasks.next().await {
710 results.push(result?);
711 }
712
713 let mut results = results.into_iter().flatten().collect::<Vec<_>>();
714
715 if let Some(limit) = limit {
716 results.truncate(limit);
717 }
718
719 let include_header = input
720 .header
721 .as_ref()
722 .and_then(|header| header.value)
723 .unwrap_or(true);
724
725 match output_format {
726 binding if binding == &OutputFormatCodes::csv() => {
727 let data = output::csv::csv(&results, include_header)?;
728
729 let base64_string: String = general_purpose::STANDARD.encode(&data);
730
731 Ok(Binary {
732 data: Some(Box::new(FHIRBase64Binary {
733 value: Some(base64_string),
734 ..Default::default()
735 })),
736 ..Default::default()
737 })
738 }
739 binding if binding == &OutputFormatCodes::json() => {
740 let data = output::json::json(&results)?;
741
742 let base64_string: String = general_purpose::STANDARD.encode(&data);
743
744 Ok(Binary {
745 data: Some(Box::new(FHIRBase64Binary {
746 value: Some(base64_string),
747 ..Default::default()
748 })),
749 ..Default::default()
750 })
751 }
752 binding if binding == &OutputFormatCodes::ndjson() => {
753 let data = output::ndjson::ndjson(results)?;
754 let base64_string: String = general_purpose::STANDARD.encode(&data);
755
756 Ok(Binary {
757 data: Some(Box::new(FHIRBase64Binary {
758 value: Some(base64_string),
759 ..Default::default()
760 })),
761 ..Default::default()
762 })
763 }
764 _ => Err(OperationOutcomeError::error(
765 IssueType::not_supported(),
766 format!("Output format '{output_format:?}' is not supported"),
767 )),
768 }
769}
770
771pub async fn view_definition_run<
783 CTX: Send + Sync + Clone + 'static,
784 Client: FHIRClient<CTX, OperationOutcomeError> + Send + Sync + 'static,
785>(
786 context: CTX,
787 client: Arc<Client>,
788 input: &ViewDefinitionRun::Input,
789) -> Result<ViewDefinitionRun::Output, OperationOutcomeError> {
790 let output_format = input
791 .format
792 .as_ref()
793 .and_then(|v| v.value.as_ref())
794 .and_then(|s| BoundCode::<OutputFormatCodes>::new(s))
795 .unwrap_or(OutputFormatCodes::csv());
796
797 let view_definition =
798 Arc::new(resolve_view_definition(context.clone(), client.as_ref(), input).await?);
799
800 let output = process_view_definition(
801 context,
802 &output_format,
803 client,
804 view_definition.as_ref(),
805 input,
806 )
807 .await?;
808
809 Ok(ViewDefinitionRun::Output { return_: output })
810}