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
104pub 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
158async 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 }
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 }
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}