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
20 changes: 20 additions & 0 deletions config/spotbugs-exclude.xml
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
<?xml version="1.0" encoding="UTF-8"?>
<FindBugsFilter>
<!-- These records defensively copy collection inputs into unmodifiable JDK collections. -->
<Match>
<Class name="~org\.offeringprotocol\.odp\.core\.(Collection|Offering.*|Page|ProblemDetails.*|SearchCapabilities.*|SearchRequests.*|ServiceDocument.*|ValidationIssue)" />
<Bug pattern="EI_EXPOSE_REP,EI_EXPOSE_REP2" />
</Match>
<Match>
<Class name="~org\.offeringprotocol\.odp\.directory\.DirectoryModels\$(Facets|PaymentFilter|SearchPage|Service|ServiceFilters|Suggestions)" />
<Bug pattern="EI_EXPOSE_REP,EI_EXPOSE_REP2" />
</Match>
<Match>
<Class name="org.offeringprotocol.odp.agent.ServiceInspection" />
<Bug pattern="EI_EXPOSE_REP,EI_EXPOSE_REP2" />
</Match>
<Match>
<Class name="~org\.offeringprotocol\.odp\.service\.OdpHttp(Request|Response)" />
<Bug pattern="EI_EXPOSE_REP,EI_EXPOSE_REP2" />
</Match>
</FindBugsFilter>
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
package org.offeringprotocol.odp.agent;

import java.net.URI;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import org.offeringprotocol.odp.core.Offering;
import org.offeringprotocol.odp.core.SearchRequests;
import org.offeringprotocol.odp.directory.DirectoryClient;
import org.offeringprotocol.odp.directory.DirectoryModels;

/** Two-stage directory-to-Service Offering discovery. */
public final class OdpAgent {
private final DirectoryClient directory;
private final ServiceClientFactory serviceClients;

public OdpAgent(DirectoryClient directory) {
this(directory, origin -> OdpServiceClient.create(URI.create(origin)));
}

public OdpAgent(DirectoryClient directory, ServiceClientFactory serviceClients) {
this.directory = Objects.requireNonNull(directory);
this.serviceClients = Objects.requireNonNull(serviceClients);
}

public List<DiscoveryEvent> searchOfferings(String query, int maximumServices, int offeringsPerService) {
if (maximumServices < 1 || maximumServices > 100) {
throw new IllegalArgumentException("maximumServices must be from 1 through 100");
}
if (offeringsPerService < 1 || offeringsPerService > 100) {
throw new IllegalArgumentException("offeringsPerService must be from 1 through 100");
}
DirectoryModels.SearchPage services =
directory.searchServices(new DirectoryModels.SearchRequest(query, null, maximumServices));
List<DiscoveryEvent> events = new ArrayList<>();
for (DirectoryModels.Service service : services.items()) {
try {
OdpServiceClient client = serviceClients.create(service.serviceOrigin());
if (!client.inspection().supports(org.offeringprotocol.odp.core.OdpOperation.SEARCH_OFFERINGS)) {
continue;
}
var request = new SearchRequests.Offerings(
org.offeringprotocol.odp.core.Odp.VERSION,
query,
null,
null,
null,
null,
null,
offeringsPerService);
for (Offering offering :
client.searchOfferings(request, "terse", null).items()) {
events.add(new OfferingEvent(service, offering));
}
} catch (RuntimeException exception) {
events.add(new IssueEvent(service, exception.getMessage()));
}
}
return List.copyOf(events);
}

@FunctionalInterface
public interface ServiceClientFactory {
OdpServiceClient create(String serviceOrigin);
}

public sealed interface DiscoveryEvent permits OfferingEvent, IssueEvent {}

public record OfferingEvent(DirectoryModels.Service service, Offering offering) implements DiscoveryEvent {}

public record IssueEvent(DirectoryModels.Service service, String message) implements DiscoveryEvent {}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
package org.offeringprotocol.odp.agent;

import java.net.http.HttpHeaders;
import org.offeringprotocol.odp.core.ProblemDetails;

public final class OdpRequestException extends RuntimeException {
private static final long serialVersionUID = 1L;

private final int responseStatus;
private final transient HttpHeaders responseHeaders;
private final transient ProblemDetails problemDetails;

public OdpRequestException(int status, String message, HttpHeaders headers, ProblemDetails problem) {
super(message);
this.responseStatus = status;
this.responseHeaders = headers;
this.problemDetails = problem;
}

public int status() {
return responseStatus;
}

public HttpHeaders headers() {
return responseHeaders;
}

public ProblemDetails problem() {
return problemDetails;
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,235 @@
package org.offeringprotocol.odp.agent;

import java.io.IOException;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.nio.charset.StandardCharsets;
import java.time.Duration;
import java.util.LinkedHashMap;
import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import org.offeringprotocol.odp.core.Collection;
import org.offeringprotocol.odp.core.Odp;
import org.offeringprotocol.odp.core.OdpJson;
import org.offeringprotocol.odp.core.OdpOperation;
import org.offeringprotocol.odp.core.OdpUris;
import org.offeringprotocol.odp.core.Offering;
import org.offeringprotocol.odp.core.OfferingPage;
import org.offeringprotocol.odp.core.OperationDescriptor;
import org.offeringprotocol.odp.core.Page;
import org.offeringprotocol.odp.core.ProblemDetails;
import org.offeringprotocol.odp.core.SearchRequests;

/** Validated ODP Service inspection and catalog client. */
public final class OdpServiceClient {
private static final int MAXIMUM_BYTES = 2_097_152;
private static final int MAXIMUM_REDIRECTS = 5;
private static final String GET = "GET";
private static final HttpClient DEFAULT_HTTP_CLIENT = HttpClient.newBuilder()
.connectTimeout(Duration.ofSeconds(10))
.followRedirects(HttpClient.Redirect.NEVER)
.build();

private final OdpTransport transport;
private final String serviceOrigin;
private final ServiceInspection serviceInspection;

private OdpServiceClient(OdpTransport transport, String serviceOrigin, ServiceInspection inspection) {
this.transport = transport;
this.serviceOrigin = serviceOrigin;
this.serviceInspection = inspection;
}

public static OdpServiceClient create(URI serviceUri) {
return create(
serviceUri, request -> DEFAULT_HTTP_CLIENT.send(request, HttpResponse.BodyHandlers.ofByteArray()));
}

public static OdpServiceClient create(URI serviceUri, OdpTransport transport) {
Objects.requireNonNull(serviceUri, "serviceUri");
Objects.requireNonNull(transport, "transport");
String origin = OdpUris.deriveServiceOrigin(serviceUri);
URI documentUri = URI.create(origin).resolve(Odp.SERVICE_DOCUMENT_PATH);
String json = request(transport, documentUri, GET, null, null, 524_288);
var document = OdpJson.parseServiceDocument(json);
Map<OdpOperation, OperationDescriptor> operations = new LinkedHashMap<>();
for (OperationDescriptor operation : document.operations()) {
operations.put(operation.name(), operation);
}
return new OdpServiceClient(
transport, origin, new ServiceInspection(origin, documentUri, document, Map.copyOf(operations)));
}

public ServiceInspection inspection() {
return serviceInspection;
}

public Page<Collection> listCollections(String representation, Integer limit, String language) {
return collectionPage(
requestOperation(OdpOperation.LIST_COLLECTIONS, null, representation, limit, language, null));
}

public Page<Collection> searchCollections(
SearchRequests.Collections request, String representation, String language) {
return collectionPage(requestOperation(
OdpOperation.SEARCH_COLLECTIONS, null, representation, null, language, OdpJson.write(request)));
}

public Collection getCollection(String id, String representation, String language) {
return OdpJson.parseCollection(
requestOperation(OdpOperation.GET_COLLECTION, id, representation, null, language, null));
}

public Page<Offering> listCollectionOfferings(
String collectionId, String representation, Integer limit, String language) {
return offeringPage(requestOperation(
OdpOperation.LIST_COLLECTION_OFFERINGS, collectionId, representation, limit, language, null));
}

public Page<Offering> listOfferings(String representation, Integer limit, String language) {
return offeringPage(requestOperation(OdpOperation.LIST_OFFERINGS, null, representation, limit, language, null));
}

public OfferingPage searchOfferings(SearchRequests.Offerings request, String representation, String language) {
OfferingPage page = OdpJson.parseOfferingSearchResponse(requestOperation(
OdpOperation.SEARCH_OFFERINGS, null, representation, null, language, OdpJson.write(request)));
page.items().forEach(item -> requireSummary(item.id(), item.name(), "Offering"));
return page;
}

public Offering getOffering(String id, String representation, String language) {
return OdpJson.parseOffering(
requestOperation(OdpOperation.GET_OFFERING, id, representation, null, language, null));
}

public Page<Collection> continueCollections(String next, String language) {
URI target = OdpUris.resolveContinuation(next, serviceOrigin);
return collectionPage(request(transport, target, GET, null, language, MAXIMUM_BYTES));
}

public Page<Offering> continueOfferings(String next, String language) {
URI target = OdpUris.resolveContinuation(next, serviceOrigin);
return offeringPage(request(transport, target, GET, null, language, MAXIMUM_BYTES));
}

private String requestOperation(
OdpOperation operation,
String identifier,
String representation,
Integer limit,
String language,
String body) {
if (!serviceInspection.supports(operation)) {
throw new IllegalStateException("Service does not advertise " + operation.value());
}
String selectedRepresentation = representation == null ? "terse" : representation;
if (!"terse".equals(selectedRepresentation) && !"full".equals(selectedRepresentation)) {
throw new IllegalArgumentException("representation must be terse or full");
}
if (limit != null && (limit < 1 || limit > 100)) {
throw new IllegalArgumentException("limit must be from 1 through 100");
}
URI target = OdpUris.buildOperationUri(
serviceInspection.document().http().endpointBase(), operation, serviceOrigin, identifier);
String separator = target.getQuery() == null ? "?" : "&";
target = URI.create(target + separator + "representation=" + selectedRepresentation
+ (limit == null ? "" : "&limit=" + limit));
return request(transport, target, operation.method(), body, language, MAXIMUM_BYTES);
}

private static Page<Collection> collectionPage(String json) {
Page<Collection> page = OdpJson.parsePage(json, Collection.class);
page.items().forEach(item -> requireSummary(item.id(), item.name(), "Collection"));
return page;
}

private static Page<Offering> offeringPage(String json) {
Page<Offering> page = OdpJson.parsePage(json, Offering.class);
page.items().forEach(item -> requireSummary(item.id(), item.name(), "Offering"));
return page;
}

private static void requireSummary(String identifier, String name, String resourceType) {
if (!OdpUris.isLocalResourceIdentifier(identifier) || name == null || name.isBlank()) {
throw new IllegalArgumentException(resourceType + " summary is invalid");
}
}

private static String request(
OdpTransport transport, URI target, String method, String body, String language, int maximumBytes) {
URI current = target;
String currentMethod = method;
String currentBody = body;
boolean hasBody = body != null;
for (int redirects = 0; redirects <= MAXIMUM_REDIRECTS; redirects++) {
HttpRequest.Builder builder = HttpRequest.newBuilder(current)
.timeout(Duration.ofSeconds(30))
.header("Accept", "application/odp+json, application/problem+json");
if (language != null && !language.isBlank()) {
builder.header("Accept-Language", language);
}
if (!hasBody) {
builder.method(currentMethod, HttpRequest.BodyPublishers.noBody());
} else {
builder.header("Content-Type", "application/odp+json")
.method(currentMethod, HttpRequest.BodyPublishers.ofString(currentBody));
}
HttpResponse<byte[]> response;
try {
response = transport.send(builder.build());
} catch (IOException exception) {
throw new IllegalStateException("ODP request failed", exception);
} catch (InterruptedException exception) {
Thread.currentThread().interrupt();
throw new IllegalStateException("ODP request was interrupted", exception);
}
int status = response.statusCode();
if (status == 301 || status == 302 || status == 303 || status == 307 || status == 308) {
if (redirects == MAXIMUM_REDIRECTS) {
throw new IllegalStateException("ODP response exceeded its redirect limit");
}
String location = response.headers()
.firstValue("Location")
.orElseThrow(() -> new IllegalStateException("ODP redirect omitted Location"));
URI next = current.resolve(location);
if (!OdpUris.deriveServiceOrigin(next).equals(OdpUris.deriveServiceOrigin(target))) {
throw new IllegalStateException("ODP redirect changed Service origin");
}
if (status == 303 || ((status == 301 || status == 302) && "POST".equals(currentMethod))) {
currentMethod = GET;
hasBody = false;
}
current = next;
} else {
byte[] bytes = response.body();
if (bytes.length > maximumBytes) {
throw new IllegalStateException("ODP response exceeds its byte limit");
}
String text = new String(bytes, StandardCharsets.UTF_8);
if (status < 200 || status > 299) {
ProblemDetails problem = null;
try {
problem = OdpJson.parseProblemDetails(text);
} catch (IllegalArgumentException ignored) {
// The HTTP status remains available when a peer does not return ODP Problem Details.
}
throw new OdpRequestException(
status,
problem == null ? "ODP request failed with HTTP " + status : problem.title(),
response.headers(),
problem);
}
String contentType =
response.headers().firstValue("Content-Type").orElse("");
if (!contentType.toLowerCase(Locale.ROOT).startsWith("application/odp+json")) {
throw new IllegalStateException("ODP response must use application/odp+json");
}
return text;
}
}
throw new IllegalStateException("ODP request produced no response");
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
package org.offeringprotocol.odp.agent;

import java.io.IOException;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;

@FunctionalInterface
public interface OdpTransport {
HttpResponse<byte[]> send(HttpRequest request) throws IOException, InterruptedException;
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
package org.offeringprotocol.odp.agent;

import java.net.URI;
import java.util.Map;
import org.offeringprotocol.odp.core.OdpOperation;
import org.offeringprotocol.odp.core.OperationDescriptor;
import org.offeringprotocol.odp.core.ServiceDocument;

public record ServiceInspection(
String serviceOrigin,
URI documentUri,
ServiceDocument document,
Map<OdpOperation, OperationDescriptor> operations) {
public ServiceInspection {
operations = Map.copyOf(operations);
}

public boolean supports(OdpOperation operation) {
return operations.containsKey(operation);
}
}
Loading