Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 6 additions & 0 deletions crates/rest/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,12 @@ helios-observability = { path = "../observability" }
# Key generation for the JWE round-trip tests
rand = "0.8"

# `SearchProvider::search_param_registry` returns a `parking_lot::RwLock`, and
# the batch unit tests' mock storage has to satisfy that trait since #478
# widened `process_batch`'s bound. Test-only: the crate's own code reaches the
# registry through `helios-persistence` and never names the lock type.
parking_lot = "0.12"

# Temp files
tempfile = "3"

Expand Down
255 changes: 246 additions & 9 deletions crates/rest/src/handlers/batch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,8 @@ use helios_audit::{AuditAction, AuditCorrelation, AuditEventBuilder};
use helios_auth::{FhirOperation, Principal, SmartScopePolicy};
use helios_fhir::FhirVersion;
use helios_persistence::core::{
BundleEntry, BundleEntryResult, BundleMethod, BundleProvider, ResourceStorage,
bundle_if_match_gate,
BundleEntry, BundleEntryResult, BundleMethod, BundleProvider, IncludeProvider, ResourceStorage,
RevincludeProvider, SearchProvider, bundle_if_match_gate,
};
use helios_persistence::error::{ResourceError, StorageError, TransactionError};
use serde_json::Value;
Expand Down Expand Up @@ -60,7 +60,13 @@ pub async fn batch_handler<S>(
request: Request,
) -> RestResult<Response>
where
S: ResourceStorage + BundleProvider + helios_persistence::core::SearchProvider + Send + Sync,
S: ResourceStorage
+ SearchProvider
+ IncludeProvider
+ RevincludeProvider
+ BundleProvider
+ Send
+ Sync,
{
// Extract the Principal from request extensions (set by auth middleware).
// If present, per-entry scope checks will be enforced.
Expand Down Expand Up @@ -251,7 +257,7 @@ async fn process_batch<S>(
principal: Option<&Principal>,
) -> RestResult<Response>
where
S: ResourceStorage + Send + Sync,
S: ResourceStorage + SearchProvider + IncludeProvider + RevincludeProvider + Send + Sync,
{
debug!(
tenant = %tenant.tenant_id(),
Expand Down Expand Up @@ -371,7 +377,13 @@ async fn process_transaction<S>(
principal: Option<&Principal>,
) -> RestResult<Response>
where
S: ResourceStorage + BundleProvider + helios_persistence::core::SearchProvider + Send + Sync,
S: ResourceStorage
+ SearchProvider
+ IncludeProvider
+ RevincludeProvider
+ BundleProvider
+ Send
+ Sync,
{
debug!(
tenant = %tenant.tenant_id(),
Expand Down Expand Up @@ -474,12 +486,43 @@ where
}
}

// GET search entries (`Patient?name=x`, bare `Patient`) cannot run inside
// the storage transaction; the spec orders GETs after all writes, so they
// execute against the just-committed state instead (#478). Their queries
// are still validated up front, where a malformed search can reject the
// whole bundle before anything executes.
let (search_entries, remaining): (Vec<_>, Vec<_>) = indexed_entries.into_iter().partition(
|(_, entry, _): &(usize, BundleEntry, Option<String>)| {
matches!(entry.method, BundleMethod::Get)
&& parse_search_entry_url(&entry.url).is_some()
},
);
let mut indexed_entries = remaining;
for (index, entry, _) in &search_entries {
let (search_type, pairs) =
parse_search_entry_url(&entry.url).expect("partitioned on is_some");
let reg = state.storage().search_param_registry(tenant.context());
let registry = reg.read();
crate::extractors::build_search_query_from_pairs(&search_type, &pairs, &registry).map_err(
|e| RestError::BadRequest {
message: format!(
"Entry {}: invalid search '{}': {}",
index,
entry.url,
e.client_response().2
),
},
)?;
}

// Conditional references (`Type?query`) resolve against the server's
// content before anything executes, per the transaction processing rules:
// exactly one match rewrites the reference to `Type/id`; zero or several
// fail the bundle (#459). They used to be stored verbatim — unsearchable
// and unresolvable. References to entries created by this same bundle use
// `fullUrl`s, which the storage layer resolves during processing.
// `fullUrl`s, which the storage layer resolves during processing. Runs on
// the write entries only — GET search entries were partitioned out above,
// and their query strings are searches, not conditional references.
resolve_conditional_references(state, &tenant, &mut indexed_entries).await?;

// Write-path validation: transactions are atomic, so any invalid write
Expand Down Expand Up @@ -539,13 +582,45 @@ where
}
}

// GET searches run against the committed state (see above). A
// failure here cannot roll the transaction back, so it surfaces
// as that entry's own error outcome rather than a misleading
// whole-bundle failure for writes that did commit.
let mut search_results: Vec<(usize, BundleEntry, BundleEntryResult)> =
Vec::with_capacity(search_entries.len());
for (index, entry, _) in &search_entries {
let (search_type, pairs) =
parse_search_entry_url(&entry.url).expect("partitioned on is_some");
let result = match crate::handlers::search::execute_search_bundle(
state,
&tenant,
&search_type,
pairs,
false,
)
.await
{
Ok(bundle) => searchset_result(bundle),
// The second of #481's two code-discarding call sites, and
// the more consequential one: this loop bypasses the
// backend executor, so it is the first *reachable*
// per-entry outcome on the transaction arm. Rendered
// through the funnel like every other entry failure.
Err(e) => entry_failure(e),
};
search_results.push((*index, entry.clone(), result));
}

// Reorder results back to original entry order
let mut ordered_results: Vec<(usize, &BundleEntry, &BundleEntryResult)> =
indexed_entries
.iter()
.zip(bundle_result.entries.iter())
.map(|((orig_idx, entry, _), result)| (*orig_idx, entry, result))
.collect();
for (orig_idx, entry, result) in &search_results {
ordered_results.push((*orig_idx, entry, result));
}
ordered_results.sort_by_key(|(idx, _, _)| *idx);

for (orig_idx, entry, result) in &ordered_results {
Expand Down Expand Up @@ -585,10 +660,15 @@ where
// rolled-back/internal transaction error never reaches the client
// response, the audit trail, or the entry outcome. The raw detail is
// preserved server-side by the `error!` log below.
// The search entries join the rollback fan-out: they were
// partitioned out of `indexed_entries` before execution, but the
// audit trail owes every entry of the failed bundle a record.
// (The status/code threading from the stacked #504 refinement
// lands with that PR; main's message-based result stays here.)
let (_, _, rollback_reason) = transaction_error_response_parts(&e);
let rollback_result =
create_error_result(500, &format!("Transaction rolled back: {rollback_reason}"));
for (orig_idx, entry, _) in &indexed_entries {
for (orig_idx, entry, _) in indexed_entries.iter().chain(&search_entries) {
let correlation_details =
EntryAuditCorrelation::from_bundle(&correlation, *orig_idx);
emit_transaction_entry_audit(
Expand Down Expand Up @@ -667,7 +747,7 @@ async fn process_batch_entry<S>(
principal: Option<&Principal>,
) -> BundleEntryResult
where
S: ResourceStorage + Send + Sync,
S: ResourceStorage + SearchProvider + IncludeProvider + RevincludeProvider + Send + Sync,
{
let request = match entry.get("request") {
Some(r) => r,
Expand Down Expand Up @@ -749,6 +829,28 @@ where

match method {
BundleMethod::Get => {
// A GET entry is either a search (`Patient?name=x`, bare
// `Patient`) or an instance read (`Patient/123`), per the spec's
// "read or search" wording for bundle GETs (#478).
if let Some((search_type, pairs)) = parse_search_entry_url(url) {
return match crate::handlers::search::execute_search_bundle(
state,
tenant,
&search_type,
pairs,
false,
)
.await
{
Ok(bundle) => searchset_result(bundle),
// Rendered through the funnel like every other entry
// failure. #481 wrote this as `let (status, _, details) =
// e.client_response()` — the same code-discard #504
// deleted everywhere else, which would have made a search
// entry the one path still answering `processing`.
Err(e) => entry_failure(e),
};
}
// Read operation
match state
.storage()
Expand Down Expand Up @@ -1316,6 +1418,54 @@ fn validation_failure_message(error: &RestError) -> String {
format!("Validation failed: {error}")
}

/// Interprets a bundle-entry GET url as a type-level search, if it is one.
///
/// Per the FHIR spec, a GET entry may carry any read OR search URL
/// (`Patient?name=x`, or bare `Patient` for an unfiltered type search).
/// Returns the resource type and the parsed query pairs, or `None` when the
/// url addresses a specific instance (`Patient/123`) and should be a read.
fn parse_search_entry_url(url: &str) -> Option<(String, Vec<(String, String)>)> {
let (path, query) = match url.split_once('?') {
Some((p, q)) => (p, Some(q)),
None => (url, None),
};
let parts: Vec<&str> = path
.trim_start_matches('/')
.split('/')
.filter(|s| !s.is_empty())
.collect();
match parts.as_slice() {
[resource_type] => Some((
resource_type.to_string(),
crate::extractors::query_pairs::parse_query_pairs(query),
)),
_ => None,
}
}

/// Builds the entry result embedding a searchset Bundle (bundle GET search).
fn searchset_result(bundle: Value) -> BundleEntryResult {
BundleEntryResult {
status: 200,
location: None,
etag: None,
last_modified: None,
resource: Some(bundle),
outcome: None,
}
}

/// Renders a failed Bundle entry through the status/details pair the
/// single-resource handlers compute.
///
/// One seam on purpose: the stacked issue-code refinement (#516) upgrades
/// this to the full `OperationOutcome` mapping; until it lands, entries keep
/// main's message-based outcomes.
fn entry_failure(err: RestError) -> BundleEntryResult {
let (status, _, details) = err.client_response();
create_error_result(status.as_u16(), &details)
}

fn create_error_result(status: u16, message: &str) -> BundleEntryResult {
let outcome = serde_json::json!({
"resourceType": "OperationOutcome",
Expand Down Expand Up @@ -1784,6 +1934,37 @@ mod tests {
use super::*;
use std::collections::{HashMap, HashSet};

/// #478: the read-or-search split on a bundle GET url. Instance addresses
/// stay reads; type-level urls (with or without a query) are searches.
#[test]
fn parse_search_entry_url_splits_reads_from_searches() {
let (rt, pairs) = parse_search_entry_url("Patient?name=x&gender=female").expect("search");
assert_eq!(rt, "Patient");
assert_eq!(pairs.len(), 2);

// Bare type: an unfiltered type search, leading slash tolerated.
let (rt, pairs) = parse_search_entry_url("/Patient").expect("bare type search");
assert_eq!(rt, "Patient");
assert!(pairs.is_empty());

// Instance reads and deeper paths are not searches.
assert!(parse_search_entry_url("Patient/123").is_none());
assert!(parse_search_entry_url("Patient/123/_history/1").is_none());
assert!(parse_search_entry_url("").is_none());
}

/// #478: the searchset entry result embeds the bundle at 200 with no
/// location/etag baggage.
#[test]
fn searchset_result_embeds_the_bundle() {
let result = searchset_result(serde_json::json!({
"resourceType": "Bundle", "type": "searchset"
}));
assert_eq!(result.status, 200);
assert!(result.location.is_none());
assert_eq!(result.resource.unwrap()["type"], "searchset");
}

use async_trait::async_trait;
use helios_audit::AuditSink;
use helios_fhir::FhirVersion;
Expand Down Expand Up @@ -2287,6 +2468,62 @@ mod tests {
}
}

// #478's search-entry dispatch widened `process_batch`'s bound to
// `SearchProvider + IncludeProvider + RevincludeProvider`, so this mock has
// to satisfy them. Every method is `unimplemented!()`, which is the same
// lever the write methods above use: no unit test in this module drives a
// search entry, and one that started to would panic loudly rather than
// silently exercising a stub.
#[async_trait]
impl helios_persistence::core::SearchProvider for DelayStorage {
async fn search(
&self,
_tenant: &TenantContext,
_query: &helios_persistence::types::SearchQuery,
) -> StorageResult<helios_persistence::core::SearchResult> {
unimplemented!()
}

async fn search_count(
&self,
_tenant: &TenantContext,
_query: &helios_persistence::types::SearchQuery,
) -> StorageResult<u64> {
unimplemented!()
}

fn search_param_registry(
&self,
_tenant: &TenantContext,
) -> Arc<parking_lot::RwLock<helios_persistence::search::SearchParameterRegistry>> {
unimplemented!()
}
}

#[async_trait]
impl helios_persistence::core::IncludeProvider for DelayStorage {
async fn resolve_includes(
&self,
_tenant: &TenantContext,
_resources: &[StoredResource],
_includes: &[helios_persistence::types::IncludeDirective],
) -> StorageResult<Vec<StoredResource>> {
unimplemented!()
}
}

#[async_trait]
impl helios_persistence::core::RevincludeProvider for DelayStorage {
async fn resolve_revincludes(
&self,
_tenant: &TenantContext,
_resources: &[StoredResource],
_revincludes: &[helios_persistence::types::IncludeDirective],
) -> StorageResult<Vec<StoredResource>> {
unimplemented!()
}
}

/// A batch Bundle of `count` GET entries, targeting `Patient/p0..p{count}`.
fn get_bundle(count: usize) -> Value {
let entries: Vec<Value> = (0..count)
Expand All @@ -2309,7 +2546,7 @@ mod tests {
principal: Option<&Principal>,
) -> Value
where
S: ResourceStorage + Send + Sync,
S: ResourceStorage + SearchProvider + IncludeProvider + RevincludeProvider + Send + Sync,
{
let tenant = TenantExtractor::new("test-tenant", crate::tenant::TenantSource::Default);
let response = process_batch(
Expand Down
Loading
Loading