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
Original file line number Diff line number Diff line change
@@ -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 {
}
69 changes: 69 additions & 0 deletions src/main/java/gov/cdc/izgateway/xform/sql/SqlFhirController.java
Original file line number Diff line number Diff line change
@@ -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<Bundle> 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<Bundle> 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<Bundle> patientMatch(
@PathVariable String name,
@PathVariable String resourceType,
HttpServletRequest req
) {
Bundle empty = new Bundle();
empty.setType(Bundle.BundleType.SEARCHSET);
return new ResponseEntity<>(empty, HttpStatus.OK);
}
}
Original file line number Diff line number Diff line change
@@ -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<OperationOutcome> 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);
}
}
Original file line number Diff line number Diff line change
@@ -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<Void> 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<Void> 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<String, Object> buildManifest(BulkExportJob job) {
return Map.of(
"transactionTime", job.getTransactionTime() != null ? job.getTransactionTime().toString() : "",
"requiresAccessToken", true,
"output", job.getOutputFiles(),
"error", java.util.List.of()
);
}
}
49 changes: 49 additions & 0 deletions src/main/java/gov/cdc/izgateway/xform/sql/bulk/BulkExportJob.java
Original file line number Diff line number Diff line change
@@ -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<OutputFile> 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<OutputFile> 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; }
}
}
Original file line number Diff line number Diff line change
@@ -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);
}
Original file line number Diff line number Diff line change
@@ -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;
}
Original file line number Diff line number Diff line change
@@ -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<UUID, BulkExportJob> 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);
}
}
Original file line number Diff line number Diff line change
@@ -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;
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
gov.cdc.izgateway.xform.sql.SqlBackendAutoConfiguration
Loading