haste_fhir_search/pg_search/
search_parameter_resolver.rs1use haste_fhir_model::r4::generated::{
2 resources::{Resource, ResourceType},
3 terminology::IssueType,
4};
5use haste_fhir_operation_error::OperationOutcomeError;
6use haste_jwt::{ProjectId, TenantId};
7use haste_repository::{Repository, fhir::CachePolicy};
8use moka::future::{Cache, CacheBuilder};
9use sqlx::{Pool, Postgres, Row};
10use std::sync::{Arc, LazyLock};
11
12use crate::{
13 ResolvedParameter, SearchParameterResolve,
14 memory::{R4_SEARCH_PARAMETERS_INDEX, SearchParametersIndex, create_index_map},
15};
16
17#[derive(Clone)]
18pub struct PgSearchParameterResolver<Repo: Repository + Send + Sync> {
19 pool: Pool<Postgres>,
20 repo: Arc<Repo>,
21}
22
23static PG_SEARCHPARAMETER_CACHE: LazyLock<
24 Cache<(TenantId, ProjectId), Arc<SearchParametersIndex>>,
25> = LazyLock::new(|| {
26 CacheBuilder::new(50_000)
27 .time_to_idle(std::time::Duration::from_hours(2))
28 .build()
29});
30
31impl<Repo: Repository + Send + Sync> PgSearchParameterResolver<Repo> {
32 pub fn new(pool: Pool<Postgres>, repo: Arc<Repo>) -> Self {
33 PgSearchParameterResolver { pool, repo }
34 }
35}
36
37async fn create_project_sp_index<Repo: Repository + Send + Sync>(
40 pool: &Pool<Postgres>,
41 repo: &Repo,
42 tenant: &TenantId,
43 project: &ProjectId,
44) -> Result<SearchParametersIndex, OperationOutcomeError> {
45 let rows = sqlx::query(
49 "SELECT sr.resource_id, sr.version_id \
50 FROM search_resource sr \
51 JOIN search_searchparameter rt ON rt.tenant = sr.tenant AND rt.project = sr.project \
52 AND rt.resource_id = sr.resource_id \
53 WHERE sr.tenant = $1 AND sr.project = $2 \
54 AND sr.resource_type = 'SearchParameter' \
55 AND 'active' = ANY(rt.status_code) \
56 LIMIT 10000",
57 )
58 .bind(tenant.as_ref())
59 .bind(project.as_ref())
60 .fetch_all(pool)
61 .await
62 .map_err(|e| {
63 OperationOutcomeError::fatal(
64 IssueType::exception(),
65 format!("Failed to query PG search for SearchParameters: {e}"),
66 )
67 })?;
68
69 let version_ids: Vec<haste_jwt::VersionId> = rows
70 .iter()
71 .map(|r| {
72 let vid: String = r.get("version_id");
73 haste_jwt::VersionId::new(vid)
74 })
75 .collect();
76
77 let version_id_refs: Vec<&haste_jwt::VersionId> = version_ids.iter().collect();
78
79 let project_sps = repo
80 .read_by_version_ids(tenant, project, &version_id_refs, CachePolicy::Cache)
81 .await?
82 .into_iter()
83 .filter_map(|r| match r {
84 Resource::SearchParameter(sp) => Some(sp),
85 _ => None,
86 })
87 .collect::<Vec<_>>();
88
89 Ok(create_index_map(
90 &crate::ParameterLevel::Project,
91 project_sps,
92 ))
93}
94
95async fn get_or_create_sp_index_for_project<Repo: Repository + Send + Sync>(
96 pool: &Pool<Postgres>,
97 repo: &Repo,
98 tenant: TenantId,
99 project: ProjectId,
100) -> Result<Option<Arc<SearchParametersIndex>>, OperationOutcomeError> {
101 if let (TenantId::System, ProjectId::System) = (&tenant, &project) {
102 return Ok(None);
103 }
104
105 let index_key = (tenant, project);
106 let pool = pool.clone();
107 let index = PG_SEARCHPARAMETER_CACHE
108 .try_get_with(index_key.clone(), async {
109 create_project_sp_index(&pool, repo, &index_key.0, &index_key.1)
110 .await
111 .map(Arc::new)
112 })
113 .await
114 .map_err(|e| OperationOutcomeError::fatal(IssueType::exception(), e.to_string()))?;
115
116 Ok(Some(index))
117}
118
119impl<Repo: Repository + Send + Sync> SearchParameterResolve for PgSearchParameterResolver<Repo> {
120 async fn by_resource_type(
121 &self,
122 tenant: &TenantId,
123 project: &ProjectId,
124 resource_type: &ResourceType,
125 ) -> Result<Vec<ResolvedParameter>, OperationOutcomeError> {
126 let mut sps = R4_SEARCH_PARAMETERS_INDEX
127 .by_resource_type(tenant, project, resource_type)
128 .await?;
129
130 if let Some(project_index) = get_or_create_sp_index_for_project(
131 &self.pool,
132 self.repo.as_ref(),
133 tenant.clone(),
134 project.clone(),
135 )
136 .await?
137 {
138 let project_sps = project_index
139 .by_resource_type(tenant, project, resource_type)
140 .await?;
141 sps.extend(project_sps);
142 }
143
144 Ok(sps)
145 }
146
147 async fn by_name(
148 &self,
149 tenant: &TenantId,
150 project: &ProjectId,
151 resource_type: Option<&ResourceType>,
152 code: &str,
153 ) -> Result<Option<ResolvedParameter>, OperationOutcomeError> {
154 if let Some(parameter) = R4_SEARCH_PARAMETERS_INDEX
155 .by_name(tenant, project, resource_type, code)
156 .await?
157 {
158 Ok(Some(parameter))
159 } else if let Some(project_index) = get_or_create_sp_index_for_project(
160 &self.pool,
161 self.repo.as_ref(),
162 tenant.clone(),
163 project.clone(),
164 )
165 .await?
166 {
167 project_index
168 .by_name(tenant, project, resource_type, code)
169 .await
170 } else {
171 Ok(None)
172 }
173 }
174
175 async fn all(
176 &self,
177 tenant: &TenantId,
178 project: &ProjectId,
179 ) -> Result<Vec<ResolvedParameter>, OperationOutcomeError> {
180 let mut all_sps = R4_SEARCH_PARAMETERS_INDEX.all(tenant, project).await?;
181
182 if let Some(project_index) = get_or_create_sp_index_for_project(
183 &self.pool,
184 self.repo.as_ref(),
185 tenant.clone(),
186 project.clone(),
187 )
188 .await?
189 {
190 all_sps.extend(project_index.all(tenant, project).await?);
191 }
192
193 Ok(all_sps)
194 }
195}