haste_fhir_search/elastic_search/
search_parameter_resolver.rs1use 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 .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}