Skip to main content

haste_fhir_search/pg_search/
search_parameter_resolver.rs

1use 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
37/// Finds active `SearchParameter` resources in the PG search index and builds
38/// a project-level search parameter index from them.
39async 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    // `SearchParameter.status` is an HL7 base parameter, so it lives in a
46    // dedicated column on the per-resource-type table rather than in the
47    // dynamic EAV tables.
48    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}