Skip to main content

Java Data Provider

This guide shows how to build a GW data provider using Spring Boot. The implementation uses dependency injection for services, a servlet filter for API-key validation, and controller classes for endpoint handlers.

Project Setup

pom.xml (key dependencies)
<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
<dependency>
<groupId>software.amazon.awssdk</groupId>
<artifactId>s3</artifactId>
</dependency>
</dependencies>
src/main/resources/application.properties
server.port=3002
gw.api-key=${GW_API_KEY:demo-api-key}
s3.endpoint=${S3_ENDPOINT:http://localstack:4566}
s3.bucket=${S3_BUCKET:gw-deliveries}

Query Envelope

Use Java records for the envelope model -- concise and immutable:

model/QueryEnvelope.java
package com.example.dataprovider.model;

import java.util.Map;

public record QueryEnvelope(
String queryId,
String datasetId,
String endpoint,
Map<String, Object> parameters,
DeliverySpec delivery,
String callbackUrl,
String callbackToken
) {
public record DeliverySpec(
String mechanism,
String url,
String s3Bucket,
String s3KeyPrefix
) {}
}

Callback Payload

model/CallbackPayload.java
package com.example.dataprovider.model;

import com.fasterxml.jackson.annotation.JsonInclude;

@JsonInclude(JsonInclude.Include.NON_NULL)
public class CallbackPayload {
private String queryId;
private String status;
private int recordCount;
private int executionTimeMs;
private DeliveryInfo delivery;
private String error;

public static class DeliveryInfo {
private String mechanism;

public DeliveryInfo() {}
public DeliveryInfo(String mechanism) { this.mechanism = mechanism; }

public String getMechanism() { return mechanism; }
public void setMechanism(String mechanism) { this.mechanism = mechanism; }
}

public CallbackPayload() {}

public CallbackPayload(String queryId, String status, int recordCount, int executionTimeMs) {
this.queryId = queryId;
this.status = status;
this.recordCount = recordCount;
this.executionTimeMs = executionTimeMs;
}

// Getters and setters
public String getQueryId() { return queryId; }
public void setQueryId(String queryId) { this.queryId = queryId; }
public String getStatus() { return status; }
public void setStatus(String status) { this.status = status; }
public int getRecordCount() { return recordCount; }
public void setRecordCount(int recordCount) { this.recordCount = recordCount; }
public int getExecutionTimeMs() { return executionTimeMs; }
public void setExecutionTimeMs(int executionTimeMs) { this.executionTimeMs = executionTimeMs; }
public DeliveryInfo getDelivery() { return delivery; }
public void setDelivery(DeliveryInfo delivery) { this.delivery = delivery; }
public String getError() { return error; }
public void setError(String error) { this.error = error; }
}

API-Key Filter

A servlet filter validates the X-GW-Api-Key header before requests reach controllers:

config/ApiKeyFilter.java
package com.example.dataprovider.config;

import jakarta.servlet.*;
import jakarta.servlet.http.HttpServletRequest;
import jakarta.servlet.http.HttpServletResponse;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.core.annotation.Order;
import org.springframework.stereotype.Component;

import java.io.IOException;

@Component
@Order(1)
public class ApiKeyFilter implements Filter {

@Value("${gw.api-key:demo-api-key}")
private String gwApiKey;

@Override
public void doFilter(ServletRequest request, ServletResponse response, FilterChain chain)
throws IOException, ServletException {

HttpServletRequest httpReq = (HttpServletRequest) request;

if ("/health".equals(httpReq.getRequestURI())) {
chain.doFilter(request, response);
return;
}

String apiKey = httpReq.getHeader("X-GW-Api-Key");
if (!gwApiKey.equals(apiKey)) {
HttpServletResponse httpRes = (HttpServletResponse) response;
httpRes.setStatus(HttpServletResponse.SC_UNAUTHORIZED);
httpRes.setContentType("application/json");
httpRes.getWriter().write("{\"error\":\"Invalid or missing X-GW-Api-Key header\"}");
return;
}

chain.doFilter(request, response);
}
}

Delivery Service

Route results to S3, a webhook, or inline:

service/DeliveryService.java
package com.example.dataprovider.service;

import com.example.dataprovider.model.QueryEnvelope;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import software.amazon.awssdk.auth.credentials.AwsBasicCredentials;
import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider;
import software.amazon.awssdk.core.sync.RequestBody;
import software.amazon.awssdk.regions.Region;
import software.amazon.awssdk.services.s3.S3Client;
import software.amazon.awssdk.services.s3.model.PutObjectRequest;

import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;

@Service
public class DeliveryService {

private final ObjectMapper objectMapper = new ObjectMapper();
private final HttpClient httpClient = HttpClient.newHttpClient();

public record DeliveryResult(String mechanism, String reference, Integer statusCode) {
public DeliveryResult(String mechanism) { this(mechanism, null, null); }
public DeliveryResult(String mechanism, String reference) { this(mechanism, reference, null); }
}

@Value("${s3.endpoint:http://localstack:4566}")
private String s3Endpoint;

@Value("${s3.bucket:gw-deliveries}")
private String s3Bucket;

public DeliveryResult deliverResults(QueryEnvelope envelope, Object results) {
var mechanism = envelope.delivery() != null ? envelope.delivery().mechanism() : "inline";
if (mechanism == null) mechanism = "inline";

return switch (mechanism) {
case "gw_s3_sync" -> deliverToS3(envelope, results);
case "webhook" -> deliverToWebhook(envelope, results);
default -> new DeliveryResult("inline");
};
}

private DeliveryResult deliverToS3(QueryEnvelope envelope, Object results) {
var spec = envelope.delivery();
var bucket = (spec.s3Bucket() != null && !spec.s3Bucket().isEmpty())
? spec.s3Bucket() : s3Bucket;
var prefix = (spec.s3KeyPrefix() != null) ? spec.s3KeyPrefix() : "";
var key = prefix + envelope.queryId() + ".json";

try {
byte[] json = objectMapper.writeValueAsBytes(results);

S3Client s3 = S3Client.builder()
.endpointOverride(URI.create(s3Endpoint))
.region(Region.US_EAST_1)
.forcePathStyle(true)
.credentialsProvider(StaticCredentialsProvider.create(
AwsBasicCredentials.create("test", "test")))
.build();

s3.putObject(
PutObjectRequest.builder()
.bucket(bucket).key(key)
.contentType("application/json")
.build(),
RequestBody.fromBytes(json));

return new DeliveryResult("gw_s3_sync", "s3://" + bucket + "/" + key);
} catch (Exception e) {
throw new RuntimeException("S3 delivery failed: " + e.getMessage(), e);
}
}

private DeliveryResult deliverToWebhook(QueryEnvelope envelope, Object results) {
var url = envelope.delivery().url();

try {
byte[] json = objectMapper.writeValueAsBytes(results);

HttpRequest request = HttpRequest.newBuilder()
.uri(URI.create(url))
.header("Content-Type", "application/json")
.POST(HttpRequest.BodyPublishers.ofByteArray(json))
.build();

HttpResponse<Void> response = httpClient.send(request, HttpResponse.BodyHandlers.discarding());
return new DeliveryResult("webhook", null, response.statusCode());
} catch (Exception e) {
throw new RuntimeException("Webhook delivery failed: " + e.getMessage(), e);
}
}
}

Callback Service

HMAC-SHA256 signing and callback posting:

service/CallbackService.java
package com.example.dataprovider.service;

import com.example.dataprovider.model.CallbackPayload;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service;

import javax.crypto.Mac;
import javax.crypto.spec.SecretKeySpec;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.nio.charset.StandardCharsets;

@Service
public class CallbackService {

private static final Logger log = LoggerFactory.getLogger(CallbackService.class);
private final ObjectMapper objectMapper = new ObjectMapper();
private final HttpClient httpClient = HttpClient.newHttpClient();

public void postCallback(String callbackUrl, String callbackToken, CallbackPayload payload) {
try {
byte[] body = objectMapper.writeValueAsBytes(payload);

// Compute HMAC-SHA256
Mac mac = Mac.getInstance("HmacSHA256");
mac.init(new SecretKeySpec(
callbackToken.getBytes(StandardCharsets.UTF_8), "HmacSHA256"));
byte[] signatureBytes = mac.doFinal(body);
String signature = bytesToHex(signatureBytes);

HttpRequest request = HttpRequest.newBuilder()
.uri(URI.create(callbackUrl))
.header("Content-Type", "application/json")
.header("X-GW-Signature", "sha256=" + signature)
.POST(HttpRequest.BodyPublishers.ofByteArray(body))
.build();

httpClient.send(request, HttpResponse.BodyHandlers.discarding());
} catch (Exception e) {
// Log but don't throw -- GW will retry or poll
log.error("[callback] Failed for query {}: {}", payload.getQueryId(), e.getMessage());
}
}

private static String bytesToHex(byte[] bytes) {
StringBuilder sb = new StringBuilder(bytes.length * 2);
for (byte b : bytes) {
sb.append(String.format("%02x", b));
}
return sb.toString();
}
}

Controller (Middleware Pattern)

The controller delegates to services for delivery and callbacks. The business logic is minimal:

controller/CompanyController.java
package com.example.dataprovider.controller;

import com.example.dataprovider.model.CallbackPayload;
import com.example.dataprovider.model.QueryEnvelope;
import com.example.dataprovider.service.CallbackService;
import com.example.dataprovider.service.DeliveryService;
import com.example.dataprovider.service.MockDataService;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.*;

import java.util.Map;

@RestController
@RequestMapping("/api/search/companies")
public class CompanyController {

private final MockDataService mockDataService;
private final CallbackService callbackService;
private final DeliveryService deliveryService;

public CompanyController(MockDataService mockDataService,
CallbackService callbackService,
DeliveryService deliveryService) {
this.mockDataService = mockDataService;
this.callbackService = callbackService;
this.deliveryService = deliveryService;
}

@PostMapping
public ResponseEntity<?> search(@RequestBody QueryEnvelope envelope) {
long startTime = System.currentTimeMillis();

if (envelope.queryId() == null || envelope.callbackUrl() == null) {
return ResponseEntity.badRequest()
.body(Map.of("error", "Missing required envelope fields"));
}

try {
// Business logic only
var params = envelope.parameters();
var name = params != null ? (String) params.get("name") : null;
var country = params != null ? (String) params.get("country") : null;
var industry = params != null ? (String) params.get("industry") : null;

var results = mockDataService.searchCompanies(name, country, industry);
int executionTimeMs = (int) (System.currentTimeMillis() - startTime);

// Deliver results
var delivery = deliveryService.deliverResults(envelope, results);

// Post success callback
var callback = new CallbackPayload(
envelope.queryId(), "delivered", results.size(), executionTimeMs);
callback.setDelivery(new CallbackPayload.DeliveryInfo(delivery.mechanism()));
callbackService.postCallback(
envelope.callbackUrl(), envelope.callbackToken(), callback);

if ("inline".equals(delivery.mechanism())) {
return ResponseEntity.ok(Map.of(
"queryId", envelope.queryId(),
"recordCount", results.size(),
"executionTimeMs", executionTimeMs,
"results", results));
}
return ResponseEntity.ok(Map.of(
"queryId", envelope.queryId(),
"recordCount", results.size(),
"executionTimeMs", executionTimeMs,
"delivery", Map.of("mechanism", delivery.mechanism(),
"reference", delivery.reference() != null ? delivery.reference() : "")));
} catch (Exception e) {
int executionTimeMs = (int) (System.currentTimeMillis() - startTime);

var errorCallback = new CallbackPayload(
envelope.queryId(), "failed", 0, executionTimeMs);
errorCallback.setError(e.getMessage());
callbackService.postCallback(
envelope.callbackUrl(), envelope.callbackToken(), errorCallback);

return ResponseEntity.internalServerError()
.body(Map.of("error", e.getMessage(), "queryId", envelope.queryId()));
}
}
}

Controller (Manual Pattern)

Explicit control over each protocol step:

controller/SanctionController.java
package com.example.dataprovider.controller;

import com.example.dataprovider.model.CallbackPayload;
import com.example.dataprovider.model.QueryEnvelope;
import com.example.dataprovider.service.CallbackService;
import com.example.dataprovider.service.DeliveryService;
import com.example.dataprovider.service.MockDataService;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.*;

import java.util.Map;

@RestController
@RequestMapping("/api/search/sanctions")
public class SanctionController {

private final MockDataService mockDataService;
private final CallbackService callbackService;
private final DeliveryService deliveryService;

public SanctionController(MockDataService mockDataService,
CallbackService callbackService,
DeliveryService deliveryService) {
this.mockDataService = mockDataService;
this.callbackService = callbackService;
this.deliveryService = deliveryService;
}

@PostMapping
public ResponseEntity<?> check(@RequestBody QueryEnvelope envelope) {
long startTime = System.currentTimeMillis();

// Step 1: Validate envelope
if (envelope.queryId() == null || envelope.callbackUrl() == null) {
return ResponseEntity.badRequest()
.body(Map.of("error", "Missing required envelope fields"));
}

// Step 2: Run business logic with custom validation
var params = envelope.parameters();
var name = params != null ? (String) params.get("name") : null;
if (name == null || name.isEmpty()) {
return ResponseEntity.badRequest()
.body(Map.of(
"error", "parameters.name is required for sanctions checks",
"queryId", envelope.queryId()));
}

var country = params != null ? (String) params.get("country") : null;
var results = mockDataService.checkSanctions(name, country);
int executionTimeMs = (int) (System.currentTimeMillis() - startTime);

try {
// Step 3: Deliver results
var delivery = deliveryService.deliverResults(envelope, results);

// Step 4: Explicitly post callback
var callback = new CallbackPayload(
envelope.queryId(), "delivered", results.size(), executionTimeMs);
callback.setDelivery(new CallbackPayload.DeliveryInfo(delivery.mechanism()));
callbackService.postCallback(
envelope.callbackUrl(), envelope.callbackToken(), callback);

if ("inline".equals(delivery.mechanism())) {
return ResponseEntity.ok(Map.of(
"queryId", envelope.queryId(),
"recordCount", results.size(),
"executionTimeMs", executionTimeMs,
"results", results));
}
return ResponseEntity.ok(Map.of(
"queryId", envelope.queryId(),
"recordCount", results.size(),
"executionTimeMs", executionTimeMs,
"delivery", Map.of("mechanism", delivery.mechanism(),
"reference", delivery.reference() != null ? delivery.reference() : "")));
} catch (Exception e) {
var errorCallback = new CallbackPayload(
envelope.queryId(), "failed", 0,
(int) (System.currentTimeMillis() - startTime));
errorCallback.setError(e.getMessage());
callbackService.postCallback(
envelope.callbackUrl(), envelope.callbackToken(), errorCallback);

return ResponseEntity.internalServerError()
.body(Map.of("error", e.getMessage(), "queryId", envelope.queryId()));
}
}
}

Running the Provider

./mvnw spring-boot:run

Testing Locally

curl -X POST http://localhost:3002/api/search/companies \
-H "Content-Type: application/json" \
-H "X-GW-Api-Key: demo-api-key" \
-d '{
"queryId": "test-001",
"datasetId": "ds-companies",
"endpoint": "/api/search/companies",
"parameters": { "name": "Acme", "country": "US" },
"delivery": { "mechanism": "inline" },
"callbackUrl": "http://localhost:8083/webhook",
"callbackToken": "test-secret"
}'