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