diff --git a/src/main/java/gov/cdc/izgateway/xform/sql/SqlBackendAutoConfiguration.java b/src/main/java/gov/cdc/izgateway/xform/sql/SqlBackendAutoConfiguration.java new file mode 100644 index 0000000..e3626b9 --- /dev/null +++ b/src/main/java/gov/cdc/izgateway/xform/sql/SqlBackendAutoConfiguration.java @@ -0,0 +1,9 @@ +package gov.cdc.izgateway.xform.sql; + +import org.springframework.context.annotation.ComponentScan; +import org.springframework.context.annotation.Configuration; + +@Configuration +@ComponentScan(basePackages = "gov.cdc.izgateway.xform.sql") +public class SqlBackendAutoConfiguration { +} diff --git a/src/main/java/gov/cdc/izgateway/xform/sql/SqlFhirController.java b/src/main/java/gov/cdc/izgateway/xform/sql/SqlFhirController.java new file mode 100644 index 0000000..7bb2231 --- /dev/null +++ b/src/main/java/gov/cdc/izgateway/xform/sql/SqlFhirController.java @@ -0,0 +1,69 @@ +package gov.cdc.izgateway.xform.sql; + +import io.swagger.v3.oas.annotations.Operation; +import io.swagger.v3.oas.annotations.responses.ApiResponse; +import jakarta.annotation.security.RolesAllowed; +import jakarta.servlet.http.HttpServletRequest; +import org.hl7.fhir.r4.model.Bundle; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.http.HttpStatus; +import org.springframework.http.ResponseEntity; +import org.springframework.web.bind.annotation.*; + +/** + * Handles single-patient FHIR queries against named SQL backends. + * Owns all paths under /sql/fhir/{name}/**, entirely distinct from + * FhirController at /fhir/**. + */ +@RestController +@RequestMapping("/sql/fhir/{name}") +@RolesAllowed({"XFORM_SENDING_SYSTEM", "ADMIN"}) +public class SqlFhirController { + + private static final Logger log = LoggerFactory.getLogger(SqlFhirController.class); + + @Operation(summary = "SQL-backed FHIR patient/immunization query") + @ApiResponse(responseCode = "200", description = "Query completed") + @GetMapping( + value = {"/{resourceType}", "/{resourceType}/_search"}, + produces = {"application/fhir+json", "application/fhir+xml", "application/json", "application/xml"} + ) + public ResponseEntity query( + @PathVariable String name, + @PathVariable String resourceType, + HttpServletRequest req + ) { + log.debug("SQL FHIR query: backend={} resource={}", name, resourceType); + Bundle empty = new Bundle(); + empty.setType(Bundle.BundleType.SEARCHSET); + empty.setTotal(0); + return new ResponseEntity<>(empty, HttpStatus.OK); + } + + @GetMapping("/{resourceType}/{id}") + public ResponseEntity read( + @PathVariable String name, + @PathVariable String resourceType, + @PathVariable String id, + HttpServletRequest req + ) { + Bundle empty = new Bundle(); + empty.setType(Bundle.BundleType.SEARCHSET); + return new ResponseEntity<>(empty, HttpStatus.OK); + } + + @PostMapping( + value = {"/{resourceType}/$match"}, + produces = {"application/fhir+json", "application/fhir+xml", "application/json"} + ) + public ResponseEntity patientMatch( + @PathVariable String name, + @PathVariable String resourceType, + HttpServletRequest req + ) { + Bundle empty = new Bundle(); + empty.setType(Bundle.BundleType.SEARCHSET); + return new ResponseEntity<>(empty, HttpStatus.OK); + } +} diff --git a/src/main/java/gov/cdc/izgateway/xform/sql/SqlUnavailableController.java b/src/main/java/gov/cdc/izgateway/xform/sql/SqlUnavailableController.java new file mode 100644 index 0000000..338a672 --- /dev/null +++ b/src/main/java/gov/cdc/izgateway/xform/sql/SqlUnavailableController.java @@ -0,0 +1,33 @@ +package gov.cdc.izgateway.xform.sql; + +import jakarta.servlet.http.HttpServletRequest; +import org.hl7.fhir.r4.model.OperationOutcome; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.context.annotation.Configuration; +import org.springframework.http.HttpStatus; +import org.springframework.http.ResponseEntity; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RestController; + +/** + * Stub controller activated when the SQL module is not configured. + * Returns 503 for all /sql/** and /bulk/sql/** paths so callers + * receive a meaningful error rather than a 404. + */ +@RestController +@ConditionalOnMissingBean(SqlBackendAutoConfiguration.class) +@Configuration +public class SqlUnavailableController { + + @RequestMapping({"/sql/**", "/bulk/sql/**"}) + public ResponseEntity unavailable(HttpServletRequest req) { + OperationOutcome oo = new OperationOutcome(); + OperationOutcome.OperationOutcomeIssueComponent issue = oo.addIssue(); + issue.setSeverity(OperationOutcome.IssueSeverity.ERROR); + issue.setCode(OperationOutcome.IssueType.NOTSUPPORTED); + issue.setDiagnostics( + "SQL backend is not configured. Deploy the izgw-transform-sql module " + + "and configure a datasource to enable this endpoint."); + return new ResponseEntity<>(oo, HttpStatus.SERVICE_UNAVAILABLE); + } +} diff --git a/src/main/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportController.java b/src/main/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportController.java new file mode 100644 index 0000000..2f93ee2 --- /dev/null +++ b/src/main/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportController.java @@ -0,0 +1,118 @@ +package gov.cdc.izgateway.xform.sql.bulk; + +import io.swagger.v3.oas.annotations.Operation; +import jakarta.annotation.security.RolesAllowed; +import jakarta.servlet.http.HttpServletRequest; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.http.HttpStatus; +import org.springframework.http.ResponseEntity; +import org.springframework.web.bind.annotation.*; + +import java.net.URI; +import java.util.Map; +import java.util.UUID; + +/** + * HL7 Bulk FHIR $export endpoints at /bulk/sql/fhir/$export. + * Follows the HL7 Bulk Data Access specification. + */ +@RestController +@RequestMapping("/bulk/sql/fhir") +@RolesAllowed({"XFORM_SENDING_SYSTEM", "BULK_EXPORT", "ADMIN"}) +public class BulkExportController { + + private static final Logger log = LoggerFactory.getLogger(BulkExportController.class); + + private final BulkExportJobStore jobStore; + private final BulkExportOutputStore outputStore; + + public BulkExportController( + @Autowired BulkExportJobStore jobStore, + @Autowired BulkExportOutputStore outputStore + ) { + this.jobStore = jobStore; + this.outputStore = outputStore; + } + + @Operation(summary = "Kick off a Bulk FHIR export job") + @PostMapping("/$export") + public ResponseEntity kickoff( + @RequestHeader(value = "Accept", required = false) String accept, + @RequestHeader(value = "Prefer", required = false) String prefer, + @RequestParam(value = "_since", required = false) String since, + @RequestParam(value = "_type", required = false) String type, + @RequestParam(value = "_typeFilter", required = false) String typeFilter, + HttpServletRequest req + ) { + if (!"respond-async".equals(prefer)) { + return ResponseEntity.badRequest().build(); + } + + if (typeFilter != null && !isSupportedTypeFilter(typeFilter)) { + return ResponseEntity.badRequest().build(); + } + + BulkExportJob job = new BulkExportJob(); + job.setSinceParam(since); + job.setTypeParam(type); + job.setTypeFilter(typeFilter); + jobStore.create(job); + + URI statusUrl = URI.create(req.getRequestURL().toString() + .replace("/$export", "/$export-status/" + job.getId())); + log.info("Bulk export job created: {}", job.getId()); + + return ResponseEntity.accepted() + .location(statusUrl) + .build(); + } + + @Operation(summary = "Poll bulk export job status") + @GetMapping("/$export-status/{jobId}") + public ResponseEntity status(@PathVariable UUID jobId) { + BulkExportJob job = jobStore.get(jobId); + if (job == null) { + return ResponseEntity.notFound().build(); + } + + return switch (job.getStatus()) { + case PENDING, RUNNING -> ResponseEntity.accepted() + .header("X-Progress", job.getStatus().name()) + .build(); + case COMPLETE -> ResponseEntity.ok(buildManifest(job)); + case FAILED -> ResponseEntity.status(HttpStatus.INTERNAL_SERVER_ERROR) + .body(Map.of("error", job.getErrorMessage())); + }; + } + + @Operation(summary = "Signal completion and schedule cleanup") + @DeleteMapping("/$export-status/{jobId}") + public ResponseEntity complete(@PathVariable UUID jobId) { + BulkExportJob job = jobStore.get(jobId); + if (job == null) { + return ResponseEntity.notFound().build(); + } + try { + outputStore.delete(jobId); + } catch (Exception e) { + log.warn("Failed to delete output files for job {}: {}", jobId, e.getMessage()); + } + jobStore.delete(jobId); + return ResponseEntity.accepted().build(); + } + + private boolean isSupportedTypeFilter(String typeFilter) { + return typeFilter.startsWith("Immunization?") || typeFilter.startsWith("Patient?"); + } + + private Map buildManifest(BulkExportJob job) { + return Map.of( + "transactionTime", job.getTransactionTime() != null ? job.getTransactionTime().toString() : "", + "requiresAccessToken", true, + "output", job.getOutputFiles(), + "error", java.util.List.of() + ); + } +} diff --git a/src/main/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportJob.java b/src/main/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportJob.java new file mode 100644 index 0000000..3518a71 --- /dev/null +++ b/src/main/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportJob.java @@ -0,0 +1,49 @@ +package gov.cdc.izgateway.xform.sql.bulk; + +import java.time.Instant; +import java.util.ArrayList; +import java.util.List; +import java.util.UUID; + +public class BulkExportJob { + + public enum Status { PENDING, RUNNING, COMPLETE, FAILED } + + private final UUID id = UUID.randomUUID(); + private volatile Status status = Status.PENDING; + private final Instant kickoffTime = Instant.now(); + private Instant transactionTime; + private String sinceParam; + private String typeFilter; + private String typeParam; + private final List outputFiles = new ArrayList<>(); + private String errorMessage; + + public UUID getId() { return id; } + public Status getStatus() { return status; } + public void setStatus(Status status) { this.status = status; } + public Instant getKickoffTime() { return kickoffTime; } + public Instant getTransactionTime() { return transactionTime; } + public void setTransactionTime(Instant transactionTime) { this.transactionTime = transactionTime; } + public String getSinceParam() { return sinceParam; } + public void setSinceParam(String sinceParam) { this.sinceParam = sinceParam; } + public String getTypeFilter() { return typeFilter; } + public void setTypeFilter(String typeFilter) { this.typeFilter = typeFilter; } + public String getTypeParam() { return typeParam; } + public void setTypeParam(String typeParam) { this.typeParam = typeParam; } + public List getOutputFiles() { return outputFiles; } + public String getErrorMessage() { return errorMessage; } + public void setErrorMessage(String errorMessage) { this.errorMessage = errorMessage; } + + public static class OutputFile { + private final String type; + private final String url; + private final int count; + public OutputFile(String type, String url, int count) { + this.type = type; this.url = url; this.count = count; + } + public String getType() { return type; } + public String getUrl() { return url; } + public int getCount() { return count; } + } +} diff --git a/src/main/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportJobStore.java b/src/main/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportJobStore.java new file mode 100644 index 0000000..1428377 --- /dev/null +++ b/src/main/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportJobStore.java @@ -0,0 +1,10 @@ +package gov.cdc.izgateway.xform.sql.bulk; + +import java.util.UUID; + +public interface BulkExportJobStore { + BulkExportJob create(BulkExportJob job); + BulkExportJob get(UUID id); + BulkExportJob update(BulkExportJob job); + void delete(UUID id); +} diff --git a/src/main/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportOutputStore.java b/src/main/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportOutputStore.java new file mode 100644 index 0000000..3d914db --- /dev/null +++ b/src/main/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportOutputStore.java @@ -0,0 +1,11 @@ +package gov.cdc.izgateway.xform.sql.bulk; + +import java.io.InputStream; +import java.io.OutputStream; +import java.util.UUID; + +public interface BulkExportOutputStore { + void write(UUID jobId, int fileIndex, InputStream data) throws Exception; + void stream(UUID jobId, int fileIndex, OutputStream out) throws Exception; + void delete(UUID jobId) throws Exception; +} diff --git a/src/main/java/gov/cdc/izgateway/xform/sql/bulk/InMemoryBulkExportJobStore.java b/src/main/java/gov/cdc/izgateway/xform/sql/bulk/InMemoryBulkExportJobStore.java new file mode 100644 index 0000000..da328d4 --- /dev/null +++ b/src/main/java/gov/cdc/izgateway/xform/sql/bulk/InMemoryBulkExportJobStore.java @@ -0,0 +1,37 @@ +package gov.cdc.izgateway.xform.sql.bulk; + +import org.springframework.stereotype.Component; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; + +/** + * V1 in-memory job store. Job state is lost on restart and not shared + * across instances. Single-instance deployment only. + */ +@Component +public class InMemoryBulkExportJobStore implements BulkExportJobStore { + + private final ConcurrentHashMap jobs = new ConcurrentHashMap<>(); + + @Override + public BulkExportJob create(BulkExportJob job) { + jobs.put(job.getId(), job); + return job; + } + + @Override + public BulkExportJob get(UUID id) { + return jobs.get(id); + } + + @Override + public BulkExportJob update(BulkExportJob job) { + jobs.put(job.getId(), job); + return job; + } + + @Override + public void delete(UUID id) { + jobs.remove(id); + } +} diff --git a/src/main/java/gov/cdc/izgateway/xform/sql/bulk/TempFileBulkExportOutputStore.java b/src/main/java/gov/cdc/izgateway/xform/sql/bulk/TempFileBulkExportOutputStore.java new file mode 100644 index 0000000..25afb1f --- /dev/null +++ b/src/main/java/gov/cdc/izgateway/xform/sql/bulk/TempFileBulkExportOutputStore.java @@ -0,0 +1,37 @@ +package gov.cdc.izgateway.xform.sql.bulk; + +import org.springframework.stereotype.Component; +import java.io.*; +import java.nio.file.*; +import java.util.UUID; + +/** + * V1 temp-file output store. Files are local to this instance and lost + * on restart. Single-instance deployment only. + */ +@Component +public class TempFileBulkExportOutputStore implements BulkExportOutputStore { + + private Path filePath(UUID jobId, int fileIndex) { + return Path.of(System.getProperty("java.io.tmpdir"), + "izg-bulk-" + jobId + "-" + fileIndex + ".ndjson"); + } + + @Override + public void write(UUID jobId, int fileIndex, InputStream data) throws Exception { + Files.copy(data, filePath(jobId, fileIndex), StandardCopyOption.REPLACE_EXISTING); + } + + @Override + public void stream(UUID jobId, int fileIndex, OutputStream out) throws Exception { + Files.copy(filePath(jobId, fileIndex), out); + } + + @Override + public void delete(UUID jobId) throws Exception { + for (int i = 0; ; i++) { + Path p = filePath(jobId, i); + if (!Files.deleteIfExists(p)) break; + } + } +} diff --git a/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports b/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports new file mode 100644 index 0000000..b41fcc2 --- /dev/null +++ b/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports @@ -0,0 +1 @@ +gov.cdc.izgateway.xform.sql.SqlBackendAutoConfiguration diff --git a/src/test/java/gov/cdc/izgateway/xform/sql/SqlFhirControllerTests.java b/src/test/java/gov/cdc/izgateway/xform/sql/SqlFhirControllerTests.java new file mode 100644 index 0000000..2f5cb65 --- /dev/null +++ b/src/test/java/gov/cdc/izgateway/xform/sql/SqlFhirControllerTests.java @@ -0,0 +1,48 @@ +package gov.cdc.izgateway.xform.sql; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.http.HttpStatus; +import org.springframework.http.ResponseEntity; +import org.springframework.mock.web.MockHttpServletRequest; + +import static org.junit.jupiter.api.Assertions.*; + +class SqlFhirControllerTests { + + private SqlFhirController controller; + + @BeforeEach + void setUp() { + controller = new SqlFhirController(); + } + + @Test + void query_returnsEmptySearchsetBundle() { + MockHttpServletRequest req = new MockHttpServletRequest("GET", "/sql/fhir/dev/Patient"); + ResponseEntity response = controller.query("dev", "Patient", req); + assertEquals(HttpStatus.OK, response.getStatusCode()); + assertNotNull(response.getBody()); + } + + @Test + void query_forImmunization_returnsOk() { + MockHttpServletRequest req = new MockHttpServletRequest("GET", "/sql/fhir/waiis/Immunization"); + ResponseEntity response = controller.query("waiis", "Immunization", req); + assertEquals(HttpStatus.OK, response.getStatusCode()); + } + + @Test + void read_returnsOk() { + MockHttpServletRequest req = new MockHttpServletRequest("GET", "/sql/fhir/dev/Patient/123"); + ResponseEntity response = controller.read("dev", "Patient", "123", req); + assertEquals(HttpStatus.OK, response.getStatusCode()); + } + + @Test + void patientMatch_returnsOk() { + MockHttpServletRequest req = new MockHttpServletRequest("POST", "/sql/fhir/dev/Patient/$match"); + ResponseEntity response = controller.patientMatch("dev", "Patient", req); + assertEquals(HttpStatus.OK, response.getStatusCode()); + } +} diff --git a/src/test/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportControllerTests.java b/src/test/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportControllerTests.java new file mode 100644 index 0000000..95ab445 --- /dev/null +++ b/src/test/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportControllerTests.java @@ -0,0 +1,77 @@ +package gov.cdc.izgateway.xform.sql.bulk; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.http.HttpStatus; +import org.springframework.http.ResponseEntity; +import org.springframework.mock.web.MockHttpServletRequest; + +import static org.junit.jupiter.api.Assertions.*; + +class BulkExportControllerTests { + + private BulkExportController controller; + private InMemoryBulkExportJobStore jobStore; + private TempFileBulkExportOutputStore outputStore; + + @BeforeEach + void setUp() { + jobStore = new InMemoryBulkExportJobStore(); + outputStore = new TempFileBulkExportOutputStore(); + controller = new BulkExportController(jobStore, outputStore); + } + + @Test + void kickoff_missingPreferHeader_returns400() { + MockHttpServletRequest req = new MockHttpServletRequest("POST", "/bulk/sql/fhir/$export"); + req.setServerName("localhost"); + req.setServerPort(443); + req.setScheme("https"); + ResponseEntity response = controller.kickoff( + "application/fhir+json", null, null, null, null, req); + assertEquals(HttpStatus.BAD_REQUEST, response.getStatusCode()); + } + + @Test + void kickoff_validRequest_returns202WithContentLocation() { + MockHttpServletRequest req = new MockHttpServletRequest("POST", "/bulk/sql/fhir/$export"); + req.setServerName("localhost"); + req.setServerPort(443); + req.setScheme("https"); + ResponseEntity response = controller.kickoff( + "application/fhir+json", "respond-async", null, null, null, req); + assertEquals(HttpStatus.ACCEPTED, response.getStatusCode()); + assertNotNull(response.getHeaders().getLocation()); + } + + @Test + void kickoff_unsupportedTypeFilter_returns400() { + MockHttpServletRequest req = new MockHttpServletRequest("POST", "/bulk/sql/fhir/$export"); + req.setServerName("localhost"); + req.setServerPort(443); + req.setScheme("https"); + ResponseEntity response = controller.kickoff( + "application/fhir+json", "respond-async", null, null, "Observation?status=final", req); + assertEquals(HttpStatus.BAD_REQUEST, response.getStatusCode()); + } + + @Test + void status_inProgress_returns202() { + BulkExportJob job = jobStore.create(new BulkExportJob()); + ResponseEntity status = controller.status(job.getId()); + assertEquals(HttpStatus.ACCEPTED, status.getStatusCode()); + } + + @Test + void delete_returns202() { + BulkExportJob job = jobStore.create(new BulkExportJob()); + ResponseEntity response = controller.complete(job.getId()); + assertEquals(HttpStatus.ACCEPTED, response.getStatusCode()); + } + + @Test + void delete_unknownJob_returns404() { + ResponseEntity response = controller.complete(java.util.UUID.randomUUID()); + assertEquals(HttpStatus.NOT_FOUND, response.getStatusCode()); + } +}