Skip to main content

haste_fhir_search/elastic_search/
search_parameter_resolver.rs

1use elasticsearch::Elasticsearch;
2use haste_fhir_client::request::{FHIRSearchTypeRequest, SearchRequest};
3use haste_fhir_model::r4::generated::resources::{Resource, ResourceType};
4use haste_fhir_operation_error::OperationOutcomeError;
5use haste_jwt::{ProjectId, TenantId};
6use haste_repository::{Repository, fhir::CachePolicy};
7use moka::future::{Cache, CacheBuilder};
8use std::sync::{Arc, LazyLock};
9
10use crate::{
11    ResolvedParameter, SearchOptions, SearchParameterResolve,
12    elastic_search::search,
13    memory::{R4_SEARCH_PARAMETERS_INDEX, SearchParametersIndex, create_index_map},
14};
15
16#[derive(Clone)]
17pub struct ElasticSearchParameterResolver<Repo: Repository + Send + Sync> {
18    es: Arc<Elasticsearch>,
19    repo: Arc<Repo>,
20}
21
22static SEARCHPARAMETER_CACHE: LazyLock<Cache<(TenantId, ProjectId), Arc<SearchParametersIndex>>> =
23    LazyLock::new(|| {
24        CacheBuilder::new(50_000)
25            // Duration for 2 hour for search parameters.
26            .time_to_idle(std::time::Duration::from_secs(2 * 60 * 60))
27            .build()
28    });
29
30impl<Repo: Repository + Send + Sync> ElasticSearchParameterResolver<Repo> {
31    pub fn new(es: Arc<Elasticsearch>, repo: Arc<Repo>) -> Self {
32        ElasticSearchParameterResolver { es, repo }
33    }
34}
35
36async fn create_project_sp_index<Repo: Repository + Send + Sync>(
37    es: Arc<Elasticsearch>,
38    repo: &Repo,
39    tenant: &TenantId,
40    project: &ProjectId,
41) -> Result<SearchParametersIndex, OperationOutcomeError> {
42    let result = search::execute_search(
43        es,
44        R4_SEARCH_PARAMETERS_INDEX.clone(),
45        &haste_repository::types::SupportedFHIRVersions::R4,
46        tenant,
47        project,
48        &SearchRequest::Type(FHIRSearchTypeRequest {
49            resource_type: ResourceType::SearchParameter,
50            parameters: vec![("status".to_string(), vec!["active".to_string()])].into(),
51        }),
52        &Some(SearchOptions {
53            count_limit: Some(10_000),
54        }),
55    )
56    .await?;
57
58    let version_ids = result
59        .entries
60        .iter()
61        .map(|r| &r.version_id)
62        .collect::<Vec<_>>();
63
64    let project_sps = repo
65        .read_by_version_ids(tenant, project, &version_ids, CachePolicy::Cache)
66        .await?
67        .into_iter()
68        .filter_map(|r| match r {
69            Resource::SearchParameter(sp) => Some(sp),
70            _ => None,
71        })
72        .collect::<Vec<_>>();
73
74    Ok(create_index_map(
75        &crate::ParameterLevel::Project,
76        project_sps,
77    ))
78}
79
80async fn get_or_create_sp_index_for_project<Repo: Repository + Send + Sync>(
81    es: Arc<Elasticsearch>,
82    repo: &Repo,
83    tenant: TenantId,
84    project: ProjectId,
85) -> Result<Option<Arc<SearchParametersIndex>>, OperationOutcomeError> {
86    match (&tenant, &project) {
87        (TenantId::System, ProjectId::System) => Ok(None),
88        _ => {
89            let index_key = (tenant, project);
90            if let Some(index) = SEARCHPARAMETER_CACHE.get(&index_key).await {
91                Ok(Some(index))
92            } else {
93                let index =
94                    Arc::new(create_project_sp_index(es, repo, &index_key.0, &index_key.1).await?);
95                SEARCHPARAMETER_CACHE.insert(index_key, index.clone()).await;
96
97                Ok(Some(index))
98            }
99        }
100    }
101}
102
103impl<Repo: Repository + Send + Sync> SearchParameterResolve
104    for ElasticSearchParameterResolver<Repo>
105{
106    async fn by_resource_type(
107        &self,
108        tenant: &haste_jwt::TenantId,
109        project: &haste_jwt::ProjectId,
110        resource_type: &haste_fhir_model::r4::generated::resources::ResourceType,
111    ) -> Result<Vec<ResolvedParameter>, OperationOutcomeError> {
112        let mut sps_by_resource_type = R4_SEARCH_PARAMETERS_INDEX
113            .by_resource_type(tenant, project, resource_type)
114            .await?;
115
116        if let Some(project_index) = get_or_create_sp_index_for_project(
117            self.es.clone(),
118            self.repo.as_ref(),
119            tenant.clone(),
120            project.clone(),
121        )
122        .await?
123        {
124            let project_sps = project_index
125                .by_resource_type(tenant, project, resource_type)
126                .await?;
127
128            sps_by_resource_type.extend(project_sps);
129        }
130
131        Ok(sps_by_resource_type)
132    }
133
134    async fn by_name(
135        &self,
136        tenant: &haste_jwt::TenantId,
137        project: &haste_jwt::ProjectId,
138        resource_type: Option<&haste_fhir_model::r4::generated::resources::ResourceType>,
139        code: &str,
140    ) -> Result<Option<ResolvedParameter>, OperationOutcomeError> {
141        if let Some(parameter) = R4_SEARCH_PARAMETERS_INDEX
142            .by_name(tenant, project, resource_type, code)
143            .await?
144        {
145            Ok(Some(parameter))
146        } else if let Some(project_index) = get_or_create_sp_index_for_project(
147            self.es.clone(),
148            self.repo.as_ref(),
149            tenant.clone(),
150            project.clone(),
151        )
152        .await?
153        {
154            project_index
155                .by_name(tenant, project, resource_type, code)
156                .await
157        } else {
158            Ok(None)
159        }
160    }
161
162    async fn all(
163        &self,
164        tenant: &haste_jwt::TenantId,
165        project: &haste_jwt::ProjectId,
166    ) -> Result<Vec<ResolvedParameter>, OperationOutcomeError> {
167        let mut all_sps = R4_SEARCH_PARAMETERS_INDEX.all(tenant, project).await?;
168
169        if let Some(project_index) = get_or_create_sp_index_for_project(
170            self.es.clone(),
171            self.repo.as_ref(),
172            tenant.clone(),
173            project.clone(),
174        )
175        .await?
176        {
177            all_sps.extend(project_index.all(tenant, project).await?);
178        }
179
180        Ok(all_sps)
181    }
182}