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