Skip to main content

haste_server/
load_artifacts.rs

1use std::{collections::HashSet, sync::Arc};
2
3use crate::{config::ServerConfig, fhir_client::ServerCTX, services::create_services};
4use base64::{Engine as _, engine::general_purpose::URL_SAFE_NO_PAD};
5use haste_artifacts::ARTIFACT_RESOURCES;
6use haste_fhir_client::{
7    FHIRClient,
8    request::{FHIRSearchTypeRequest, SearchRequest},
9    url::ParsedParameter,
10};
11use haste_fhir_model::r4::generated::{
12    resources::{Resource, ResourceType, SearchParameter, StructureDefinition},
13    terminology::IssueType,
14    types::{Coding, FHIRCode, FHIRUri, Meta},
15};
16use haste_fhir_operation_error::OperationOutcomeError;
17use haste_fhir_search::{SearchEngine, SearchOptions};
18use haste_jwt::{ProjectId, TenantId};
19
20use haste_repository::{Repository, fhir::CachePolicy, types::SupportedFHIRVersions};
21use sha1::{Digest, Sha1};
22
23fn generate_sha256_hash(value: &Resource) -> String {
24    let json = serde_json::to_string(value).expect("failed to serialize value.");
25    let mut sha_hasher = Sha1::new();
26    sha_hasher.update(json.as_bytes());
27    let sha1 = sha_hasher.finalize();
28
29    URL_SAFE_NO_PAD.encode(sha1)
30}
31
32static HASH_TAG_SYSTEM: &str = "https://haste.health/fhir/CodeSystem/hash";
33
34fn _add_hash_tag(meta: &mut Option<Box<Meta>>, sha_hash: String) {
35    let hash_tag = Coding {
36        system: Some(Box::new(FHIRUri {
37            value: Some(HASH_TAG_SYSTEM.to_string()),
38            ..Default::default()
39        })),
40        code: Some(Box::new(FHIRCode {
41            value: Some(sha_hash),
42            ..Default::default()
43        })),
44        ..Default::default()
45    };
46
47    let meta = if let Some(meta) = meta {
48        meta
49    } else {
50        *meta = Some(Box::new(Meta::default()));
51        meta.as_mut().unwrap()
52    };
53
54    match &mut meta.tag {
55        Some(tags) => tags.push(hash_tag),
56        None => meta.tag = Some(vec![hash_tag]),
57    }
58}
59
60fn add_hash_tag(resource: &mut Resource, sha_hash: String) {
61    match resource {
62        Resource::StructureDefinition(structure_definition) => {
63            _add_hash_tag(&mut structure_definition.meta, sha_hash)
64        }
65        Resource::CodeSystem(code_system) => _add_hash_tag(&mut code_system.meta, sha_hash),
66        Resource::ValueSet(value_set) => _add_hash_tag(&mut value_set.meta, sha_hash),
67        Resource::SearchParameter(search_parameter) => {
68            _add_hash_tag(&mut search_parameter.meta, sha_hash)
69        }
70        _ => {}
71    }
72}
73
74fn get_id(resource: &Resource) -> String {
75    match resource {
76        Resource::StructureDefinition(structure_definition) => {
77            structure_definition.id.clone().unwrap_or_default()
78        }
79        Resource::CodeSystem(code_system) => code_system.id.clone().unwrap_or_default(),
80        Resource::ValueSet(value_set) => value_set.id.clone().unwrap_or_default(),
81        Resource::SearchParameter(search_parameter) => {
82            search_parameter.id.clone().unwrap_or_default()
83        }
84        _ => todo!(
85            "Unsupported resource type '{}'",
86            resource.resource_type().as_ref()
87        ),
88    }
89}
90
91pub fn get_resource_type(resource: &Resource) -> ResourceType {
92    match resource {
93        Resource::StructureDefinition(_) => ResourceType::StructureDefinition,
94        Resource::CodeSystem(_) => ResourceType::CodeSystem,
95        Resource::ValueSet(_) => ResourceType::ValueSet,
96        Resource::SearchParameter(_) => ResourceType::SearchParameter,
97        _ => todo!(
98            "Unsupported resource type '{}'",
99            resource.resource_type().as_ref()
100        ),
101    }
102}
103
104/// This deletes existing artifacts and then reloads them. In a single transaction.
105pub async fn reset_artifacts(config: Arc<ServerConfig>) -> Result<(), OperationOutcomeError> {
106    let services = create_services(config.clone()).await?;
107
108    let transaction = services.transaction().await?;
109
110    {
111        let ctx = Arc::new(ServerCTX::system(
112            TenantId::System,
113            ProjectId::System,
114            transaction.fhir_client.clone(),
115            transaction.rate_limit.clone(),
116        ));
117
118        tracing::info!("Deleting existing CodeSystems");
119        ctx.client
120            .delete_type(
121                ctx.clone(),
122                ResourceType::CodeSystem,
123                (vec![] as Vec<(String, Vec<String>)>).into(),
124            )
125            .await?;
126        tracing::info!("Deleting existing ValueSets");
127        ctx.client
128            .delete_type(
129                ctx.clone(),
130                ResourceType::ValueSet,
131                (vec![] as Vec<(String, Vec<String>)>).into(),
132            )
133            .await?;
134        tracing::info!("Deleting existing StructureDefinitions");
135        ctx.client
136            .delete_type(
137                ctx.clone(),
138                ResourceType::StructureDefinition,
139                (vec![] as Vec<(String, Vec<String>)>).into(),
140            )
141            .await?;
142        tracing::info!("Deleting existing SearchParameters");
143        ctx.client
144            .delete_type(
145                ctx.clone(),
146                ResourceType::SearchParameter,
147                (vec![] as Vec<(String, Vec<String>)>).into(),
148            )
149            .await?;
150        _load_artifacts(ctx.clone()).await?;
151    }
152
153    transaction.commit().await?;
154
155    Ok(())
156}
157
158// Used for both reloading artifacts and reset.
159async fn _load_artifacts<Client: FHIRClient<Arc<ServerCTX<Client>>, OperationOutcomeError>>(
160    ctx: Arc<ServerCTX<Client>>,
161) -> Result<(), OperationOutcomeError> {
162    let mut hashes = HashSet::new();
163
164    let mut total_loaded: usize = 0;
165    for resource in ARTIFACT_RESOURCES.iter() {
166        let sha_hash = generate_sha256_hash(resource);
167        hashes.insert(sha_hash);
168
169        match &resource {
170            Resource::SearchParameter(_)
171            | Resource::CodeSystem(_)
172            | Resource::ValueSet(_)
173            | Resource::StructureDefinition(_) => {
174                let mut resource = resource.clone();
175                let resource_type = get_resource_type(&resource);
176                let id = get_id(&resource);
177                let sha_hash = generate_sha256_hash(&resource);
178
179                add_hash_tag(&mut resource, sha_hash.clone());
180
181                let res = ctx
182                    .client
183                    .conditional_update(
184                        ctx.clone(),
185                        resource_type.clone(),
186                        vec![
187                            ParsedParameter::from(("_id".to_string(), vec![id.clone()])),
188                            ParsedParameter::from((
189                                "_tag".to_string(),
190                                vec![HASH_TAG_SYSTEM.to_string() + "|" + sha_hash.as_str()],
191                                "not".to_string(),
192                            )),
193                        ]
194                        .into(),
195                        resource.clone(),
196                    )
197                    .await;
198
199                if let Ok(res) = res {
200                    total_loaded += 1;
201                    tracing::info!(
202                        "Updated '{}' with id '{}' and sha '{}'",
203                        resource_type.as_ref(),
204                        res.id().as_deref().unwrap_or("unknown"),
205                        sha_hash.as_str()
206                    );
207                } else if let Err(err) = res {
208                    let code = &err.outcome().issue[0].code;
209                    let diagnostic = err.outcome().issue[0]
210                        .diagnostics
211                        .as_deref()
212                        .and_then(|d| d.value.as_deref())
213                        .unwrap_or("unknown");
214
215                    match &err.outcome().issue[0].code {
216                        i if i == &IssueType::invalid() => {
217                            tracing::error!("{:#?}", err);
218                            panic!("INVALID");
219                        }
220                        i if i == &IssueType::conflict() => {
221                            // Ignore.
222                        }
223                        _ => {
224                            tracing::error!(
225                                "Failed to update '{}' with id '{}'. Issue code: '{:?}', diagnostic: '{}'",
226                                resource_type.as_ref(),
227                                id,
228                                code,
229                                diagnostic
230                            );
231                        }
232                    }
233                }
234            }
235            _ => {
236                // println!("Skipping resource.");
237            }
238        }
239    }
240
241    tracing::info!(
242        "Loaded a total of '{}' artifacts with unique hashes '{}'",
243        total_loaded,
244        hashes.len(),
245    );
246
247    Ok(())
248}
249
250pub async fn load_artifacts(config: Arc<ServerConfig>) -> Result<(), OperationOutcomeError> {
251    let services = create_services(config.clone()).await?;
252
253    let ctx = Arc::new(ServerCTX::system(
254        TenantId::System,
255        ProjectId::System,
256        services.fhir_client.clone(),
257        services.rate_limit.clone(),
258    ));
259
260    _load_artifacts(ctx.clone()).await
261}
262
263pub async fn get_all_sds<Repo: Repository, Search: SearchEngine>(
264    kinds: &[&str],
265    repo: &Repo,
266    search_engine: &Search,
267) -> Result<Vec<StructureDefinition>, OperationOutcomeError> {
268    let sd_search = FHIRSearchTypeRequest {
269        resource_type: ResourceType::StructureDefinition,
270        parameters: vec![
271            (
272                "kind".to_string(),
273                kinds.iter().map(|s| s.to_string()).collect(),
274            ),
275            ("abstract".to_string(), vec!["false".to_string()]),
276            ("derivation".to_string(), vec!["specialization".to_string()]),
277        ]
278        .into(),
279    };
280    let sd_results = search_engine
281        .search(
282            &SupportedFHIRVersions::R4,
283            &TenantId::System,
284            &ProjectId::System,
285            &SearchRequest::Type(sd_search),
286            Some(SearchOptions {
287                count_limit: Some(10_000),
288            }),
289        )
290        .await?;
291
292    let version_ids = sd_results
293        .entries
294        .iter()
295        .map(|v| &v.version_id)
296        .collect::<Vec<_>>();
297
298    let sds = repo
299        .read_by_version_ids(
300            &TenantId::System,
301            &ProjectId::System,
302            version_ids.as_slice(),
303            CachePolicy::NoCache,
304        )
305        .await?
306        .into_iter()
307        .filter_map(|r| match r {
308            Resource::StructureDefinition(sd) => Some(sd),
309            _ => None,
310        });
311
312    Ok(sds.collect())
313}
314
315pub async fn get_all_sps<Repo: Repository, Search: SearchEngine>(
316    repo: &Repo,
317    search_engine: &Search,
318) -> Result<Vec<SearchParameter>, OperationOutcomeError> {
319    let sp_search = FHIRSearchTypeRequest {
320        resource_type: ResourceType::SearchParameter,
321        parameters: (vec![] as Vec<(String, Vec<String>)>).into(),
322    };
323    let sp_results = search_engine
324        .search(
325            &SupportedFHIRVersions::R4,
326            &TenantId::System,
327            &ProjectId::System,
328            &SearchRequest::Type(sp_search),
329            Some(SearchOptions {
330                count_limit: Some(10_000),
331            }),
332        )
333        .await?;
334
335    let version_ids = sp_results
336        .entries
337        .iter()
338        .map(|v| &v.version_id)
339        .collect::<Vec<_>>();
340
341    let sps = repo
342        .read_by_version_ids(
343            &TenantId::System,
344            &ProjectId::System,
345            version_ids.as_slice(),
346            CachePolicy::NoCache,
347        )
348        .await?
349        .into_iter()
350        .filter_map(|r| match r {
351            Resource::SearchParameter(sp) => Some(sp),
352            _ => None,
353        });
354
355    Ok(sps.collect())
356}