getAsStringAsync() {
+ return getBytesAsync().thenApply(b -> new String(b, charset != null ? charset : StandardCharsets.UTF_8));
}
public int getCode() {
diff --git a/Plan/api/src/main/java/com/djrapitops/plan/delivery/web/resolver/ResponseBuilder.java b/Plan/api/src/main/java/com/djrapitops/plan/delivery/web/resolver/ResponseBuilder.java
index 00311f7485..b5c49eba82 100644
--- a/Plan/api/src/main/java/com/djrapitops/plan/delivery/web/resolver/ResponseBuilder.java
+++ b/Plan/api/src/main/java/com/djrapitops/plan/delivery/web/resolver/ResponseBuilder.java
@@ -21,6 +21,7 @@
import java.nio.charset.Charset;
import java.nio.charset.StandardCharsets;
+import java.util.concurrent.CompletableFuture;
public class ResponseBuilder {
@@ -30,17 +31,6 @@ public class ResponseBuilder {
this.response = new Response();
}
- /**
- * Set MIME Type of the Response.
- *
- * @param mimeType MIME type of the Response Documentation
- * @return this builder.
- * @see MimeType for common MIME types.
- */
- public ResponseBuilder setMimeType(String mimeType) {
- return setHeader("Content-Type", mimeType);
- }
-
/**
* Set HTTP Status code.
*
@@ -86,11 +76,17 @@ public ResponseBuilder setContent(WebResource resource) {
}
public ResponseBuilder setContent(byte[] bytes) {
- response.bytes = bytes;
- return setHeader("Content-Length", bytes.length)
+ byte[] safeBytes = bytes != null ? bytes : new byte[0];
+ response.bytes = CompletableFuture.completedFuture(safeBytes);
+ return setHeader("Content-Length", safeBytes.length)
.setHeader("Accept-Ranges", "bytes"); // Does not compress
}
+ public ResponseBuilder setContent(CompletableFuture bytesFuture) {
+ response.bytes = bytesFuture != null ? bytesFuture : CompletableFuture.completedFuture(new byte[0]);
+ return this;
+ }
+
public ResponseBuilder setContent(String utf8String) {
return setContent(utf8String, StandardCharsets.UTF_8);
}
@@ -112,6 +108,22 @@ public ResponseBuilder setContent(String content, Charset charset) {
.removeHeader("Accept-Ranges"); // Can compress
}
+ public ResponseBuilder setContent(CompletableFuture stringFuture, Charset charset) {
+ if (stringFuture == null) return setContent(new byte[0]);
+ Charset effectiveCharset = charset != null ? charset : StandardCharsets.UTF_8;
+ String mimeType = getMimeType();
+ response.charset = effectiveCharset;
+
+ if (mimeType != null) {
+ String[] parts = mimeType.split(";");
+ if (parts.length == 1) {
+ setMimeType(parts[0] + "; charset=" + effectiveCharset.name().toLowerCase());
+ }
+ }
+ removeHeader("Accept-Ranges");
+ return setContent(stringFuture.thenApply(string -> string == null ? null : string.getBytes(effectiveCharset)));
+ }
+
/**
* Set content as serialized JSON object.
*
@@ -120,6 +132,10 @@ public ResponseBuilder setContent(String content, Charset charset) {
*/
public ResponseBuilder setJSONContent(Object objectToSerialize) {
if (objectToSerialize instanceof String) return setJSONContent((String) objectToSerialize);
+ if (objectToSerialize instanceof CompletableFuture) {
+ CompletableFuture> future = (CompletableFuture>) objectToSerialize;
+ return setJSONContent(future.thenApply(obj -> obj instanceof String ? (String) obj : new Gson().toJson(obj)));
+ }
return setJSONContent(new Gson().toJson(objectToSerialize));
}
@@ -127,6 +143,10 @@ public ResponseBuilder setJSONContent(String json) {
return setMimeType(MimeType.JSON).setContent(json);
}
+ public ResponseBuilder setJSONContent(CompletableFuture jsonFuture) {
+ return setMimeType(MimeType.JSON).setContent(jsonFuture, StandardCharsets.UTF_8);
+ }
+
/**
* Finish building.
*
@@ -137,15 +157,16 @@ public ResponseBuilder setJSONContent(String json) {
* @see #setMimeType(String) to set MIME-type.
*/
public Response build() {
- byte[] content = response.bytes;
- if(content == null && response.code == 204) {
+ CompletableFuture contentFuture = response.bytes;
+ if (contentFuture == null && response.code == 204) {
// HTTP Code 204 requires no response, so there is no need to validate it.
return response;
}
- exceptionIf(content == null, "Content not defined for Response");
+ exceptionIf(contentFuture == null, "Content not defined for Response");
String mimeType = getMimeType();
- exceptionIf(content.length > 0 && mimeType == null, "MIME Type not defined for Response");
- exceptionIf(content.length > 0 && mimeType.isEmpty(), "MIME Type empty for Response");
+ boolean hasContent = response.bytes != null && response.bytes.isDone() && response.bytes.join().length > 0;
+ exceptionIf(hasContent && mimeType == null, "MIME Type not defined for Response");
+ exceptionIf(hasContent && mimeType.isEmpty(), "MIME Type empty for Response");
exceptionIf(response.code < 100 || response.code >= 600, "HTTP Status code out of bounds (" + response.code + ")");
return response;
}
@@ -154,6 +175,17 @@ private String getMimeType() {
return response.headers.get("Content-Type");
}
+ /**
+ * Set MIME Type of the Response.
+ *
+ * @param mimeType MIME type of the Response Documentation
+ * @return this builder.
+ * @see MimeType for common MIME types.
+ */
+ public ResponseBuilder setMimeType(String mimeType) {
+ return setHeader("Content-Type", mimeType);
+ }
+
private void exceptionIf(boolean value, String errorMsg) {
if (value) throw new InvalidResponseException(errorMsg);
}
diff --git a/Plan/api/src/main/java/com/djrapitops/plan/delivery/web/resolver/request/Request.java b/Plan/api/src/main/java/com/djrapitops/plan/delivery/web/resolver/request/Request.java
index c8e1c08d53..42496a63ec 100644
--- a/Plan/api/src/main/java/com/djrapitops/plan/delivery/web/resolver/request/Request.java
+++ b/Plan/api/src/main/java/com/djrapitops/plan/delivery/web/resolver/request/Request.java
@@ -21,6 +21,7 @@
import java.util.Map;
import java.util.Optional;
+import java.util.concurrent.CompletableFuture;
/**
* Represents a HTTP request to use with {@link Resolver}.
@@ -34,7 +35,7 @@ public final class Request {
private final URIQuery query;
private final WebUser user;
private final Map headers;
- private final byte[] requestBody;
+ private final CompletableFuture requestBody;
private final String accessIpAddress;
/**
@@ -65,12 +66,27 @@ public Request(String method, URIPath path, URIQuery query, WebUser user, Map headers, byte[] requestBody, String accessIpAddress) {
+ this(method, path, query, user, headers, CompletableFuture.completedFuture(requestBody != null ? requestBody : new byte[0]), accessIpAddress);
+ }
+
+ /**
+ * Constructor.
+ *
+ * @param method HTTP method, GET, PUT, POST, etc
+ * @param path Requested path /example/target
+ * @param query Request parameters ?param=value etc
+ * @param user Web user doing the request (if authenticated)
+ * @param headers Request headers Documentation
+ * @param requestBody CompletableFuture of raw body bytes
+ * @param accessIpAddress IP address this request is coming from.
+ */
+ public Request(String method, URIPath path, URIQuery query, WebUser user, Map headers, CompletableFuture requestBody, String accessIpAddress) {
this.method = method;
this.path = path;
this.query = query;
this.user = user;
this.headers = headers;
- this.requestBody = requestBody;
+ this.requestBody = requestBody != null ? requestBody : CompletableFuture.completedFuture(new byte[0]);
this.accessIpAddress = accessIpAddress;
}
@@ -109,7 +125,7 @@ public Request(String method, String target, WebUser user, Map h
}
this.user = user;
this.headers = headers;
- this.requestBody = new byte[0];
+ this.requestBody = CompletableFuture.completedFuture(new byte[0]);
this.accessIpAddress = accessIpAddress;
}
@@ -142,10 +158,20 @@ public URIQuery getQuery() {
/**
* Get the raw body, if present.
+ * Blocks until the body is available if fetched asynchronously.
*
* @return byte[].
*/
public byte[] getRequestBody() {
+ return requestBody.join();
+ }
+
+ /**
+ * Get the raw body as a {@link CompletableFuture}.
+ *
+ * @return CompletableFuture of byte[].
+ */
+ public CompletableFuture getRequestBodyAsync() {
return requestBody;
}
@@ -184,7 +210,7 @@ public String toString() {
", query=" + query +
", user=" + user +
", headers=" + headers +
- ", body=" + requestBody.length +
+ ", body=" + (requestBody.isDone() && !requestBody.isCompletedExceptionally() ? requestBody.join().length : "async") +
'}';
}
}
diff --git a/Plan/build.gradle b/Plan/build.gradle
index 0183b88e46..772fe46d62 100644
--- a/Plan/build.gradle
+++ b/Plan/build.gradle
@@ -23,7 +23,7 @@ clean {
allprojects {
ext {
majorVersion = "5"
- minorVersion = "8"
+ minorVersion = "9"
buildVersion = providers.provider {
def command = "git rev-list --count HEAD"
def buildInfo = command.execute().text.trim()
@@ -37,7 +37,7 @@ allprojects {
}
group = "com.djrapitops"
- version = project.hasProperty("isRelease") ? "$fullVersionFilename" : "5.8-SNAPSHOT"
+ version = project.hasProperty("isRelease") ? "$fullVersionFilename" : "5.9-SNAPSHOT"
}
subprojects {
@@ -69,7 +69,7 @@ subprojects {
commonsCodecVersion = "1.22.1"
caffeineVersion = "3.2.4"
jetbrainsAnnotationsVersion = "26.1.0"
- jettyVersion = "11.0.26"
+ jettyVersion = "12.1.12"
mysqlVersion = "9.7.0"
mariadbVersion = "3.5.10"
sqliteVersion = "3.42.0.1"
diff --git a/Plan/bukkit/build.gradle b/Plan/bukkit/build.gradle
index 992fe0d338..5e5f080492 100644
--- a/Plan/bukkit/build.gradle
+++ b/Plan/bukkit/build.gradle
@@ -21,7 +21,7 @@ dependencies {
}
compileJava {
- options.release = 11
+ options.release = 17
}
processResources {
diff --git a/Plan/bungeecord/build.gradle b/Plan/bungeecord/build.gradle
index a09ec9da85..83c53c8961 100644
--- a/Plan/bungeecord/build.gradle
+++ b/Plan/bungeecord/build.gradle
@@ -14,7 +14,7 @@ dependencies {
}
compileJava {
- options.release = 11
+ options.release = 17
}
processResources {
diff --git a/Plan/common/build.gradle b/Plan/common/build.gradle
index 4c6e47e6e7..727b3c01eb 100644
--- a/Plan/common/build.gradle
+++ b/Plan/common/build.gradle
@@ -81,7 +81,7 @@ dependencies {
api "com.google.code.gson:gson:$gsonVersion"
api "org.eclipse.jetty:jetty-server:$jettyVersion"
implementation "org.eclipse.jetty:jetty-alpn-java-server:$jettyVersion"
- implementation "org.eclipse.jetty.http2:http2-server:$jettyVersion"
+ implementation "org.eclipse.jetty.http2:jetty-http2-server:$jettyVersion"
implementation "org.jasypt:jasypt:$jasyptVersion:lite"
// Swagger annotations
@@ -119,7 +119,7 @@ dependencies {
}
compileJava {
- options.release = 11
+ options.release = 17
}
test {
diff --git a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/RequestBodyConverter.java b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/RequestBodyConverter.java
index 6e14165b9b..e97057403e 100644
--- a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/RequestBodyConverter.java
+++ b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/RequestBodyConverter.java
@@ -36,14 +36,26 @@ private RequestBodyConverter() {
* @return {@link URIQuery}.
*/
public static URIQuery formBody(Request request) {
- return new URIQuery(new String(request.getRequestBody(), StandardCharsets.UTF_8));
+ return formBody(request.getRequestBody());
+ }
+
+ public static URIQuery formBody(byte[] bytes) {
+ return new URIQuery(new String(bytes != null ? bytes : new byte[0], StandardCharsets.UTF_8));
}
public static T bodyJson(@Untrusted Request request, Gson gson, Class ofType) {
- return gson.fromJson(new String(request.getRequestBody(), StandardCharsets.UTF_8), ofType);
+ return bodyJson(request.getRequestBody(), gson, ofType);
+ }
+
+ public static T bodyJson(@Untrusted byte[] bytes, Gson gson, Class ofType) {
+ return gson.fromJson(new String(bytes != null ? bytes : new byte[0], StandardCharsets.UTF_8), ofType);
}
public static T bodyJson(Request request, Gson gson, TypeToken ofType) {
- return gson.fromJson(new String(request.getRequestBody(), StandardCharsets.UTF_8), ofType.getType());
+ return bodyJson(request.getRequestBody(), gson, ofType);
+ }
+
+ public static T bodyJson(byte[] bytes, Gson gson, TypeToken ofType) {
+ return gson.fromJson(new String(bytes != null ? bytes : new byte[0], StandardCharsets.UTF_8), ofType.getType());
}
}
diff --git a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/ResponseResolver.java b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/ResponseResolver.java
index abe51356a6..c32a9bcf4b 100644
--- a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/ResponseResolver.java
+++ b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/ResponseResolver.java
@@ -18,12 +18,12 @@
import com.djrapitops.plan.delivery.web.ResolverService;
import com.djrapitops.plan.delivery.web.ResolverSvc;
+import com.djrapitops.plan.delivery.web.resolver.AsyncResolver;
import com.djrapitops.plan.delivery.web.resolver.NoAuthResolver;
import com.djrapitops.plan.delivery.web.resolver.Resolver;
import com.djrapitops.plan.delivery.web.resolver.Response;
import com.djrapitops.plan.delivery.web.resolver.exception.BadRequestException;
import com.djrapitops.plan.delivery.web.resolver.exception.MethodNotAllowedException;
-import com.djrapitops.plan.delivery.web.resolver.exception.NotFoundException;
import com.djrapitops.plan.delivery.web.resolver.request.Request;
import com.djrapitops.plan.delivery.web.resolver.request.WebUser;
import com.djrapitops.plan.delivery.webserver.auth.FailReason;
@@ -36,7 +36,6 @@
import com.djrapitops.plan.delivery.webserver.resolver.swagger.SwaggerPageResolver;
import com.djrapitops.plan.exceptions.WebUserAuthException;
import com.djrapitops.plan.utilities.dev.Untrusted;
-import com.djrapitops.plan.utilities.logging.ErrorContext;
import com.djrapitops.plan.utilities.logging.ErrorLogger;
import dagger.Lazy;
import io.swagger.v3.oas.annotations.OpenAPIDefinition;
@@ -46,8 +45,11 @@
import javax.inject.Inject;
import javax.inject.Singleton;
+import java.util.Iterator;
import java.util.List;
import java.util.Optional;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionException;
import java.util.function.Supplier;
import java.util.regex.Pattern;
@@ -188,64 +190,87 @@ private NoAuthResolver fileResolver(Supplier response) {
return request -> Optional.of(response.get());
}
- public Response getResponse(@Untrusted Request request) {
+ public CompletableFuture getResponse(@Untrusted Request request) {
try {
- return tryToGetResponse(request);
- } catch (BadRequestException e) {
- return responseFactory.badRequest(e.getMessage(), request.getPath().asString());
- } catch (NotFoundException e) {
- return responseFactory.notFound404(e.getMessage());
- } catch (MethodNotAllowedException e) {
- return responseFactory.methodNotAllowed405(e.getMessage(), e.getAllowedMethods());
- } catch (WebUserAuthException e) {
- throw e; // Pass along
- } catch (Exception e) {
- errorLogger.error(e, ErrorContext.builder().related(request).build());
- return responseFactory.internalErrorResponse(e, "Failed to get a response");
+ return tryToGetResponse(request)
+ .exceptionallyCompose(throwable -> handleException(request, throwable));
+ } catch (Exception t) {
+ return handleException(request, t);
}
}
- /**
- * @throws NotFoundException In some cases when page was not found, not all.
- * @throws BadRequestException If the request did not have required things.
- */
- private Response tryToGetResponse(@Untrusted Request request) {
+ public CompletableFuture handleException(@Untrusted Request request, Throwable exception) {
+ while (exception instanceof CompletionException || exception instanceof java.util.concurrent.ExecutionException) {
+ if (exception.getCause() != null) {
+ exception = exception.getCause();
+ } else {
+ break;
+ }
+ }
+ if (exception instanceof BadRequestException badRequest) {
+ return CompletableFuture.completedFuture(responseFactory.badRequest(
+ badRequest.getMessage(), request.getPath().asString()));
+ }
+ if (exception instanceof MethodNotAllowedException notAllowed) {
+ return CompletableFuture.completedFuture(responseFactory.methodNotAllowed405(
+ notAllowed.getMessage(), notAllowed.getAllowedMethods()));
+ }
+ if (exception instanceof WebUserAuthException authException) {
+ return CompletableFuture.failedFuture(authException); // Pass along
+ }
+ return CompletableFuture.completedFuture(responseFactory.internalErrorResponse(
+ exception, "Failed to get a response"));
+ }
+
+ private CompletableFuture tryToGetResponse(@Untrusted Request request) {
if ("OPTIONS".equalsIgnoreCase(request.getMethod())) {
// https://developer.mozilla.org/en-US/docs/Web/HTTP/Methods/OPTIONS
- return Response.builder().setStatus(204).build();
+ return CompletableFuture.completedFuture(Response.builder().setStatus(204).build());
}
- Optional user = request.getUser();
List foundResolvers = resolverService.getResolvers(request.getPath().asString());
- if (foundResolvers.isEmpty()) return responseFactory.pageNotFound404();
-
- for (Resolver resolver : foundResolvers) {
- boolean isAuthRequired = webServer.get().isAuthRequired() && resolver.requiresAuth(request);
- if (isAuthRequired) {
- if (user.isEmpty()) {
- if (webServer.get().isUsingHTTPS()) {
- throw new WebUserAuthException(FailReason.NO_USER_PRESENT);
- } else {
- return responseFactory.forbidden403();
- }
+ if (foundResolvers.isEmpty()) return CompletableFuture.completedFuture(responseFactory.pageNotFound404());
+ return resolveNext(foundResolvers.iterator(), request);
+ }
+
+ private CompletableFuture resolveNext(Iterator iterator, @Untrusted Request request) {
+ if (!iterator.hasNext()) {
+ return CompletableFuture.completedFuture(responseFactory.pageNotFound404());
+ }
+
+ Resolver resolver = iterator.next();
+ Optional user = request.getUser();
+ boolean isAuthRequired = webServer.get().isAuthRequired() && resolver.requiresAuth(request);
+ if (isAuthRequired) {
+ if (user.isEmpty()) {
+ if (webServer.get().isUsingHTTPS()) {
+ throw new WebUserAuthException(FailReason.NO_USER_PRESENT);
+ } else {
+ return CompletableFuture.completedFuture(responseFactory.forbidden403());
}
+ }
- if (resolver.canAccess(request)) {
- Optional resolved = resolver.resolve(request);
- if (resolved.isPresent()) return resolved.get();
+ if (!resolver.canAccess(request)) {
+ if (request.getPath().startsWith("/v1/")) {
+ return CompletableFuture.completedFuture(responseFactory.forbidden403Json());
} else {
- if (request.getPath().startsWith("/v1/")) {
- return responseFactory.forbidden403Json();
- } else {
- return responseFactory.forbidden403();
- }
+ return CompletableFuture.completedFuture(responseFactory.forbidden403());
}
- } else {
- Optional resolved = resolver.resolve(request);
- if (resolved.isPresent()) return resolved.get();
}
}
- return responseFactory.pageNotFound404();
+
+ CompletableFuture> futureResponse;
+ try {
+ futureResponse = resolver instanceof AsyncResolver asyncResolver
+ ? asyncResolver.resolveAsync(request)
+ : CompletableFuture.supplyAsync(() -> resolver.resolve(request));
+ } catch (Exception t) { // Some exceptions may be thrown sync based on the request.
+ futureResponse = CompletableFuture.failedFuture(t);
+ }
+ return futureResponse.thenCompose(gotResponse ->
+ gotResponse.map(CompletableFuture::completedFuture) // try next resolver
+ .orElseGet(() -> resolveNext(iterator, request))
+ );
}
}
diff --git a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyInternalRequest.java b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyInternalRequest.java
index 5cd00b0e21..b9c3645ae9 100644
--- a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyInternalRequest.java
+++ b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyInternalRequest.java
@@ -23,63 +23,61 @@
import com.djrapitops.plan.delivery.webserver.auth.Cookie;
import com.djrapitops.plan.delivery.webserver.configuration.WebserverConfiguration;
import com.djrapitops.plan.utilities.dev.Untrusted;
-import jakarta.servlet.http.HttpServletRequest;
-import org.apache.commons.text.TextStringBuilder;
+import org.eclipse.jetty.http.HttpCookie;
+import org.eclipse.jetty.http.HttpField;
import org.eclipse.jetty.http.HttpHeader;
+import org.eclipse.jetty.io.Content;
import org.eclipse.jetty.server.Request;
+import org.eclipse.jetty.util.Promise;
-import java.io.BufferedReader;
-import java.io.ByteArrayOutputStream;
-import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
-import java.util.Spliterators;
-import java.util.function.Function;
+import java.util.concurrent.CompletableFuture;
import java.util.stream.Collectors;
-import java.util.stream.Stream;
-import java.util.stream.StreamSupport;
public class JettyInternalRequest implements InternalRequest {
- private final Request baseRequest;
- private final HttpServletRequest request;
+ private final Request request;
private final WebserverConfiguration webserverConfiguration;
private final AuthenticationExtractor authenticationExtractor;
- public JettyInternalRequest(Request baseRequest, HttpServletRequest request, WebserverConfiguration webserverConfiguration, AuthenticationExtractor authenticationExtractor) {
- this.baseRequest = baseRequest;
+ public JettyInternalRequest(Request request, WebserverConfiguration webserverConfiguration, AuthenticationExtractor authenticationExtractor) {
this.request = request;
this.webserverConfiguration = webserverConfiguration;
this.authenticationExtractor = authenticationExtractor;
}
@Override
- public long getTimestamp() {return baseRequest.getTimeStamp();}
+ public long getTimestamp() {return Request.getTimeStamp(request);}
@Override
- public String getMethod() {return baseRequest.getMethod();}
+ public String getMethod() {return request.getMethod();}
@Override
public String getAccessAddressFromSocketIp() {
- return baseRequest.getRemoteAddr();
+ return Request.getRemoteAddr(request);
}
@Override
public String getAccessAddressFromHeader() {
- String header = baseRequest.getHeader(HttpHeader.X_FORWARDED_FOR.asString());
+ String header = getHeader(HttpHeader.X_FORWARDED_FOR);
if (header != null && header.contains(",")) {
return header.split(",")[0].trim();
}
return header;
}
+ private String getHeader(HttpHeader headerName) {
+ return request.getHeaders().get(headerName);
+ }
+
@Override
public com.djrapitops.plan.delivery.web.resolver.request.Request toRequest(@Untrusted String accessAddress) {
- String requestMethod = baseRequest.getMethod();
- @Untrusted URIPath path = new URIPath(baseRequest.getHttpURI().getDecodedPath());
- @Untrusted URIQuery query = new URIQuery(baseRequest.getHttpURI().getQuery());
- @Untrusted byte[] requestBody = readRequestBody();
+ String requestMethod = request.getMethod();
+ @Untrusted URIPath path = new URIPath(request.getHttpURI().getDecodedPath());
+ @Untrusted URIQuery query = new URIQuery(request.getHttpURI().getQuery());
+ CompletableFuture requestBody = readRequestBodyAsync();
WebUser user = getWebUser(webserverConfiguration, authenticationExtractor, accessAddress);
@Untrusted Map headers = getRequestHeaders();
return new com.djrapitops.plan.delivery.web.resolver.request.Request(requestMethod, path, query, user, headers, requestBody, accessAddress);
@@ -87,62 +85,53 @@ public com.djrapitops.plan.delivery.web.resolver.request.Request toRequest(@Untr
@Override
public Map getRequestHeaders() {
- return streamHeaderNames()
- .collect(Collectors.toMap(Function.identity(), baseRequest::getHeader,
+ return request.getHeaders().stream()
+ .collect(Collectors.toMap(HttpField::getName, HttpField::getValue,
(one, two) -> one + ';' + two));
}
- private Stream streamHeaderNames() {
- return StreamSupport.stream(Spliterators.spliteratorUnknownSize(baseRequest.getHeaderNames().asIterator(), 0), false);
- }
+ private CompletableFuture readRequestBodyAsync() {
+ CompletableFuture future = new CompletableFuture<>();
+ Content.Source.asByteArrayAsync(request, Integer.MAX_VALUE, new Promise.Invocable<>() {
+ @Override
+ public void succeeded(byte[] result) {
+ future.complete(result != null ? result : new byte[0]);
+ }
- private byte[] readRequestBody() {
- try (BufferedReader reader = request.getReader();
- ByteArrayOutputStream buf = new ByteArrayOutputStream(512)) {
- int b;
- while ((b = reader.read()) != -1) {
- buf.write((byte) b);
+ @Override
+ public void failed(Throwable x) {
+ future.complete(new byte[0]);
}
- return buf.toByteArray();
- } catch (IOException ignored) {
- // requestBody stays empty
- return new byte[0];
- }
+ });
+ return future;
}
@Override
public List getCookies() {
- @Untrusted List textCookies = getCookieHeaders();
+ @Untrusted List jettyCookies = Request.getCookies(request);
List cookies = new ArrayList<>();
- if (!textCookies.isEmpty()) {
- String[] separated = new TextStringBuilder().appendWithSeparators(textCookies, ";").get().split(";");
- for (String textCookie : separated) {
- cookies.add(new Cookie(textCookie.trim()));
+ if (!jettyCookies.isEmpty()) {
+ for (HttpCookie cookie : jettyCookies) {
+ cookies.add(new Cookie(cookie.getName(), cookie.getValue()));
}
}
return cookies;
}
- private List getCookieHeaders() {
- return StreamSupport.stream(Spliterators.spliteratorUnknownSize(request.getHeaders(HttpHeader.COOKIE.asString()).asIterator(), 0), false)
- .collect(Collectors.toList());
- }
-
@Override
public String getRequestedURIString() {
- return baseRequest.getRequestURI();
+ return request.getHttpURI().getPath();
}
@Override
public String getRequestedPathAndQuery() {
- return baseRequest.getHttpURI().getDecodedPath() + baseRequest.getHttpURI().getQuery();
+ return request.getHttpURI().getDecodedPath() + request.getHttpURI().getQuery();
}
@Override
public String toString() {
return "JettyInternalRequest{" +
- "baseRequest=" + baseRequest +
- ", request=" + request +
+ "request=" + request +
", webserverConfiguration=" + webserverConfiguration +
", authenticationExtractor=" + authenticationExtractor +
'}';
diff --git a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyRequestHandler.java b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyRequestHandler.java
index f1186721c7..36fd3ef064 100644
--- a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyRequestHandler.java
+++ b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyRequestHandler.java
@@ -16,63 +16,79 @@
*/
package com.djrapitops.plan.delivery.webserver.http;
-import com.djrapitops.plan.delivery.web.resolver.Response;
import com.djrapitops.plan.delivery.webserver.Addresses;
import com.djrapitops.plan.delivery.webserver.auth.AuthenticationExtractor;
import com.djrapitops.plan.delivery.webserver.configuration.WebserverConfiguration;
+import com.djrapitops.plan.processing.Processing;
import com.djrapitops.plan.settings.config.PlanConfig;
import com.djrapitops.plan.settings.config.paths.PluginSettings;
import com.djrapitops.plan.utilities.logging.ErrorContext;
import com.djrapitops.plan.utilities.logging.ErrorLogger;
-import jakarta.servlet.ServletException;
-import jakarta.servlet.http.HttpServletRequest;
-import jakarta.servlet.http.HttpServletResponse;
import net.playeranalytics.plugin.server.PluginLogger;
+import org.eclipse.jetty.server.Handler;
import org.eclipse.jetty.server.Request;
-import org.eclipse.jetty.server.handler.AbstractHandler;
+import org.eclipse.jetty.util.Callback;
import javax.inject.Inject;
import javax.inject.Singleton;
-import java.io.IOException;
+import java.util.concurrent.CompletableFuture;
+import java.util.function.Function;
@Singleton
-public class JettyRequestHandler extends AbstractHandler {
+public class JettyRequestHandler extends Handler.Abstract {
private final WebserverConfiguration webserverConfiguration;
private final AuthenticationExtractor authenticationExtractor;
private final Addresses addresses;
private final RequestHandler requestHandler;
+ private final Processing processing;
private final PlanConfig config;
private final PluginLogger logger;
private final ErrorLogger errorLogger;
@Inject
- public JettyRequestHandler(WebserverConfiguration webserverConfiguration, AuthenticationExtractor authenticationExtractor, Addresses addresses, RequestHandler requestHandler, PlanConfig config, PluginLogger logger, ErrorLogger errorLogger) {
+ public JettyRequestHandler(WebserverConfiguration webserverConfiguration, AuthenticationExtractor authenticationExtractor, Addresses addresses, RequestHandler requestHandler, Processing processing, PlanConfig config, PluginLogger logger, ErrorLogger errorLogger) {
this.webserverConfiguration = webserverConfiguration;
this.authenticationExtractor = authenticationExtractor;
this.addresses = addresses;
this.requestHandler = requestHandler;
+ this.processing = processing;
this.config = config;
this.logger = logger;
this.errorLogger = errorLogger;
}
@Override
- public void handle(String target, Request baseRequest, HttpServletRequest servletRequest, HttpServletResponse servletResponse) throws IOException, ServletException {
+ public boolean handle(Request request, org.eclipse.jetty.server.Response jettyResponse, Callback callback) throws Exception {
try {
- InternalRequest internalRequest = new JettyInternalRequest(baseRequest, servletRequest, webserverConfiguration, authenticationExtractor);
- Response response = requestHandler.getResponse(internalRequest);
- new JettyResponseSender(response, servletRequest, servletResponse, addresses).send();
- baseRequest.setHandled(true);
+ InternalRequest internalRequest = new JettyInternalRequest(request, webserverConfiguration, authenticationExtractor);
+ CompletableFuture.supplyAsync(() -> requestHandler.getResponse(internalRequest), processing.getNonCriticalExecutor())
+ .thenCompose(Function.identity())
+ .thenApply(response -> new JettyResponseSender(response, request, jettyResponse, addresses))
+ .thenCompose(JettyResponseSender::sendAsync)
+ .whenComplete((result, throwable) -> {
+ if (throwable != null) {
+ logError(request, throwable);
+ callback.failed(throwable);
+ } else {
+ callback.succeeded();
+ }
+ });
+ return true;
} catch (Exception e) {
- if (config.isTrue(PluginSettings.DEV_MODE)) {
- logger.warn("THIS ERROR IS ONLY LOGGED IN DEV MODE:");
- errorLogger.warn(e, ErrorContext.builder()
- .whatToDo("THIS ERROR IS ONLY LOGGED IN DEV MODE")
- .related(baseRequest.getMethod(), baseRequest.getRemoteAddr(), target, baseRequest.getRequestURI())
- .build());
- }
+ logError(request, e);
+ callback.failed(e);
+ throw e;
}
+ }
+ private void logError(Request request, Throwable e) {
+ if (config.isTrue(PluginSettings.DEV_MODE)) {
+ logger.warn("THIS ERROR IS ONLY LOGGED IN DEV MODE:");
+ errorLogger.warn(e, ErrorContext.builder()
+ .whatToDo("THIS ERROR IS ONLY LOGGED IN DEV MODE")
+ .related(request.getMethod(), Request.getRemoteAddr(request), request.getHttpURI().getPath())
+ .build());
+ }
}
}
diff --git a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyResponseSender.java b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyResponseSender.java
index d05f2a90ea..b6ec231380 100644
--- a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyResponseSender.java
+++ b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyResponseSender.java
@@ -17,59 +17,70 @@
package com.djrapitops.plan.delivery.webserver.http;
import com.djrapitops.plan.delivery.web.resolver.MimeType;
-import com.djrapitops.plan.delivery.web.resolver.Response;
import com.djrapitops.plan.delivery.webserver.Addresses;
-import jakarta.servlet.http.HttpServletRequest;
-import jakarta.servlet.http.HttpServletResponse;
import org.apache.commons.lang3.Strings;
import org.eclipse.jetty.http.HttpHeader;
+import org.eclipse.jetty.server.Request;
+import org.eclipse.jetty.server.Response;
+import org.eclipse.jetty.util.Callback;
-import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
-import java.io.OutputStream;
+import java.nio.ByteBuffer;
import java.util.Map;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionException;
import java.util.zip.GZIPOutputStream;
public class JettyResponseSender {
- private final Response response;
- private final HttpServletRequest servletRequest;
- private final HttpServletResponse servletResponse;
+ private final com.djrapitops.plan.delivery.web.resolver.Response response;
+ private final Request jettyRequest;
+ private final Response jettyResponse;
private final Addresses addresses;
- public JettyResponseSender(Response response, HttpServletRequest servletRequest, HttpServletResponse servletResponse, Addresses addresses) {
+ public JettyResponseSender(com.djrapitops.plan.delivery.web.resolver.Response response, Request jettyRequest, Response jettyResponse, Addresses addresses) {
this.response = response;
- this.servletRequest = servletRequest;
- this.servletResponse = servletResponse;
+ this.jettyRequest = jettyRequest;
+ this.jettyResponse = jettyResponse;
this.addresses = addresses;
}
public void send() throws IOException {
- if ("HEAD".equals(servletRequest.getMethod()) || response.getCode() == 204 || response.getCode() == 304) {
+ try {
+ sendAsync().join();
+ } catch (CompletionException e) {
+ if (e.getCause() instanceof IOException) {
+ throw (IOException) e.getCause();
+ }
+ throw new IOException(e.getCause() != null ? e.getCause() : e);
+ }
+ }
+
+ public CompletableFuture sendAsync() {
+ if ("HEAD".equals(jettyRequest.getMethod()) || response.getCode() == 204 || response.getCode() == 304) {
setResponseHeaders();
- sendHeadResponse();
+ return sendHeadResponse();
} else if (canGzip()) {
- sendCompressed();
+ return sendCompressed();
} else {
setResponseHeaders();
- sendRawBytes();
+ return sendRawBytes();
}
}
private boolean canGzip() {
- String method = servletRequest.getMethod();
+ String method = jettyRequest.getMethod();
String mimeType = response.getHeaders().get(HttpHeader.CONTENT_TYPE.asString());
return "GET".equals(method) && Strings.CS.containsAny(mimeType, MimeType.HTML, MimeType.CSS, MimeType.JS, MimeType.JSON, "text/plain");
}
- public void sendHeadResponse() throws IOException {
- try {
- response.getHeaders().remove(HttpHeader.CONTENT_LENGTH.asString());
- beginSend();
- } finally {
- servletResponse.getOutputStream().close();
- }
+ public CompletableFuture sendHeadResponse() {
+ response.getHeaders().remove(HttpHeader.CONTENT_LENGTH.asString());
+ beginSend();
+ Callback.Completable completable = new Callback.Completable();
+ jettyResponse.write(true, ByteBuffer.wrap(new byte[0]), completable);
+ return completable;
}
private void setResponseHeaders() {
@@ -77,7 +88,7 @@ private void setResponseHeaders() {
correctRedirect(responseHeaders);
for (Map.Entry header : responseHeaders.entrySet()) {
- servletResponse.setHeader(header.getKey(), header.getValue());
+ jettyResponse.getHeaders().add(header.getKey(), header.getValue());
}
}
@@ -89,26 +100,34 @@ private void correctRedirect(Map responseHeaders) {
}
}
- private void sendCompressed() throws IOException {
- response.getHeaders().remove(HttpHeader.ACCEPT_RANGES.asString());
- response.getHeaders().put(HttpHeader.CONTENT_ENCODING.asString(), "gzip");
-
- byte[] gzipped = gzip();
- try (OutputStream out = servletResponse.getOutputStream()) {
- response.getHeaders().put(HttpHeader.CONTENT_LENGTH.asString(), String.valueOf(gzipped.length));
- setResponseHeaders();
-
- servletResponse.setStatus(response.getCode());
-
- send(out, gzipped);
- }
+ private CompletableFuture sendCompressed() {
+ return response.getBytesAsync().thenCompose(rawBytes -> {
+ try {
+ response.getHeaders().remove(HttpHeader.ACCEPT_RANGES.asString());
+ response.getHeaders().put(HttpHeader.CONTENT_ENCODING.asString(), "gzip");
+
+ byte[] gzipped = gzip(rawBytes);
+ response.getHeaders().put(HttpHeader.CONTENT_LENGTH.asString(), String.valueOf(gzipped.length));
+ setResponseHeaders();
+
+ jettyResponse.setStatus(response.getCode());
+
+ Callback.Completable completable = new Callback.Completable();
+ jettyResponse.write(true, ByteBuffer.wrap(gzipped), completable);
+ return completable;
+ } catch (IOException e) {
+ CompletableFuture failed = new CompletableFuture<>();
+ failed.completeExceptionally(e);
+ return failed;
+ }
+ });
}
- private byte[] gzip() throws IOException {
+ private byte[] gzip(byte[] bytes) throws IOException {
try (ByteArrayOutputStream bufferStream = new ByteArrayOutputStream();
GZIPOutputStream gzipStream = new GZIPOutputStream(bufferStream)
) {
- gzipStream.write(response.getBytes());
+ gzipStream.write(bytes);
gzipStream.finish();
gzipStream.flush();
return bufferStream.toByteArray();
@@ -121,35 +140,21 @@ private void beginSend() {
|| "0".equals(length)
|| response.getCode() == 204
|| response.getCode() == 304
- || "HEAD".equals(servletRequest.getMethod())
+ || "HEAD".equals(jettyRequest.getMethod())
) {
- servletResponse.setHeader(HttpHeader.CONTENT_LENGTH.asString(), null);
+ jettyResponse.getHeaders().remove(HttpHeader.CONTENT_LENGTH.asString());
}
// Return a content length of -1 for HTTP code 204 (No content)
// and HEAD requests to avoid warning messages.
- servletResponse.setStatus(response.getCode());
+ jettyResponse.setStatus(response.getCode());
}
- private void sendRawBytes() throws IOException {
+ private CompletableFuture sendRawBytes() {
beginSend();
- try (OutputStream out = servletResponse.getOutputStream()) {
- send(out);
- }
- }
-
- private void send(OutputStream out) throws IOException {
- send(out, response.getBytes());
- }
-
- private void send(OutputStream out, byte[] bytes) throws IOException {
- try (
- ByteArrayInputStream bis = new ByteArrayInputStream(bytes)
- ) {
- byte[] buffer = new byte[2048];
- int count;
- while ((count = bis.read(buffer)) != -1) {
- out.write(buffer, 0, count);
- }
- }
+ return response.getBytesAsync().thenCompose(bytes -> {
+ Callback.Completable completable = new Callback.Completable();
+ jettyResponse.write(true, ByteBuffer.wrap(bytes), completable);
+ return completable;
+ });
}
}
diff --git a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyWebserver.java b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyWebserver.java
index 8af80aee67..54b31f6106 100644
--- a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyWebserver.java
+++ b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/JettyWebserver.java
@@ -76,76 +76,80 @@ public void enable() {
return;
}
- webserver = new Server();
- webserver.setStopAtShutdown(true);
-
- this.port = webserverConfiguration.getPort();
-
- HttpConfiguration configuration = new HttpConfiguration();
- Optional sslContext = getSslContextFactory();
- sslContext.ifPresent(ssl -> {
- configuration.setSecureScheme("https");
- configuration.setSecurePort(port);
-
- SecureRequestCustomizer serverNameIdentifierCheckSkipper = new SecureRequestCustomizer();
- serverNameIdentifierCheckSkipper.setSniHostCheck(false);
- serverNameIdentifierCheckSkipper.setSniRequired(false);
- configuration.addCustomizer(serverNameIdentifierCheckSkipper);
-
- usingHttps = true;
- });
-
- HttpConnectionFactory httpConnector = new HttpConnectionFactory(configuration);
-
- HTTP2CServerConnectionFactory http2CConnector = new HTTP2CServerConnectionFactory(configuration);
- http2CConnector.setConnectProtocolEnabled(true);
-
-
- ServerConnector connector = sslContext
- .map(sslContextFactory -> {
- HTTP2ServerConnectionFactory http2Connector = new HTTP2ServerConnectionFactory(configuration);
- http2Connector.setConnectProtocolEnabled(true);
- ALPNServerConnectionFactory alpn = getAlpnServerConnectionFactory(httpConnector.getProtocol());
-
- return new ServerConnector(webserver, sslContextFactory, alpn, httpConnector, http2Connector, http2CConnector);
- })
- .orElseGet(() -> {
- if (webserverConfiguration.isProxyModeHttps()) {
- webserverLogMessages.authenticationUsingProxy();
- } else {
- webserverLogMessages.authenticationNotPossible();
- }
- return new ServerConnector(webserver, httpConnector, http2CConnector);
- });
-
- connector.setPort(port);
- String internalIP = webserverConfiguration.getInternalIP();
- connector.setHost(internalIP);
- webserver.addConnector(connector);
-
- webserver.setHandler(jettyRequestHandler);
-
- String startFailure = "Failed to start Jetty webserver: ";
- try {
- webserver.start();
- } catch (IOException e) {
- if (e.getMessage().contains("Failed to bind")) {
- boolean defaultInternalIp = "0.0.0.0".equals(internalIP);
- String causeHelp = defaultInternalIp ? ", is the port (" + port + ") in use?" : ", is the Internal_IP (" + internalIP + ") invalid? (Use 0.0.0.0 for automatic)";
- throw new EnableException(startFailure + e.getMessage().replace("0.0.0.0", "") + causeHelp, e);
- } else {
+ // Jetty loads its services from thread context classloader, but thread context doesn't have all Jetty classes.
+ // Because of this we swap to plugin classloader before any operations so that the loading succeeds.
+ ThreadContextClassLoaderSwap.performOperation(getClass().getClassLoader(), () -> {
+ webserver = new Server();
+ webserver.setStopAtShutdown(true);
+
+ this.port = webserverConfiguration.getPort();
+
+ HttpConfiguration configuration = new HttpConfiguration();
+ Optional sslContext = getSslContextFactory();
+ sslContext.ifPresent(ssl -> {
+ configuration.setSecureScheme("https");
+ configuration.setSecurePort(port);
+
+ SecureRequestCustomizer serverNameIdentifierCheckSkipper = new SecureRequestCustomizer();
+ serverNameIdentifierCheckSkipper.setSniHostCheck(false);
+ serverNameIdentifierCheckSkipper.setSniRequired(false);
+ configuration.addCustomizer(serverNameIdentifierCheckSkipper);
+
+ usingHttps = true;
+ });
+
+ HttpConnectionFactory httpConnector = new HttpConnectionFactory(configuration);
+
+ HTTP2CServerConnectionFactory http2CConnector = new HTTP2CServerConnectionFactory(configuration);
+ http2CConnector.setConnectProtocolEnabled(true);
+
+
+ ServerConnector connector = sslContext
+ .map(sslContextFactory -> {
+ HTTP2ServerConnectionFactory http2Connector = new HTTP2ServerConnectionFactory(configuration);
+ http2Connector.setConnectProtocolEnabled(true);
+ ALPNServerConnectionFactory alpn = new ALPNServerConnectionFactory("h2", "h2c", "http/1.1");
+ alpn.setDefaultProtocol(httpConnector.getProtocol());
+ return new ServerConnector(webserver, sslContextFactory, alpn, httpConnector, http2Connector, http2CConnector);
+ })
+ .orElseGet(() -> {
+ if (webserverConfiguration.isProxyModeHttps()) {
+ webserverLogMessages.authenticationUsingProxy();
+ } else {
+ webserverLogMessages.authenticationNotPossible();
+ }
+ return new ServerConnector(webserver, httpConnector, http2CConnector);
+ });
+
+ connector.setPort(port);
+ String internalIP = webserverConfiguration.getInternalIP();
+ connector.setHost(internalIP);
+ webserver.addConnector(connector);
+
+ webserver.setHandler(jettyRequestHandler);
+
+ String startFailure = "Failed to start Jetty webserver: ";
+ try {
+ webserver.start();
+ } catch (IOException e) {
+ if (e.getMessage().contains("Failed to bind")) {
+ boolean defaultInternalIp = "0.0.0.0".equals(internalIP);
+ String causeHelp = defaultInternalIp ? ", is the port (" + port + ") in use?" : ", is the Internal_IP (" + internalIP + ") invalid? (Use 0.0.0.0 for automatic)";
+ throw new EnableException(startFailure + e.getMessage().replace("0.0.0.0", "") + causeHelp, e);
+ } else {
+ throw new EnableException(startFailure + e.toString(), e);
+ }
+ } catch (Exception e) {
throw new EnableException(startFailure + e.toString(), e);
}
- } catch (Exception e) {
- throw new EnableException(startFailure + e.toString(), e);
- }
- webserverLogMessages.infoWebserverEnabled(getPort());
- sslContext.map(SslContextFactory::getKeyStore).ifPresent(this::logCertificateExpiryInformation);
+ webserverLogMessages.infoWebserverEnabled(getPort());
+ sslContext.map(SslContextFactory::getKeyStore).ifPresent(this::logCertificateExpiryInformation);
- responseResolver.registerPages();
+ responseResolver.registerPages();
- webserverConfiguration.getAllowedIpList().prepare();
+ webserverConfiguration.getAllowedIpList().prepare();
+ });
}
private void logCertificateExpiryInformation(KeyStore keyStore) {
@@ -164,22 +168,6 @@ private void logCertificateExpiryInformation(KeyStore keyStore) {
}
}
- private ALPNServerConnectionFactory getAlpnServerConnectionFactory(String protocol) {
- ClassLoader pluginClassLoader = getClass().getClassLoader();
- return ThreadContextClassLoaderSwap.performOperation(pluginClassLoader, () -> {
- try {
- Class.forName("org.eclipse.jetty.alpn.java.server.JDK9ServerALPNProcessor");
- // ALPN is protocol upgrade protocol required for upgrading http 1.1 connections to 2
- ALPNServerConnectionFactory alpn = new ALPNServerConnectionFactory("h2", "h2c", "http/1.1");
- alpn.setDefaultProtocol(protocol);
- return alpn;
- } catch (IllegalStateException | ClassNotFoundException ignored) {
- logger.warn("JDK9ServerALPNProcessor not found. ALPN (HTTP/2 upgrade protocol) is not available.");
- return null;
- }
- });
- }
-
private Optional getSslContextFactory() {
if (webserverConfiguration.isProxyModeHttps()) {
return Optional.empty();
diff --git a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/RequestHandler.java b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/RequestHandler.java
index 6b9b6c868a..f1f6402943 100644
--- a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/RequestHandler.java
+++ b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/http/RequestHandler.java
@@ -31,7 +31,9 @@
import javax.inject.Inject;
import javax.inject.Singleton;
+import java.util.Map;
import java.util.Optional;
+import java.util.concurrent.CompletableFuture;
@Singleton
public class RequestHandler {
@@ -55,61 +57,87 @@ public RequestHandler(WebserverConfiguration webserverConfiguration, ResponseFac
rateLimitGuard = new RateLimitGuard();
}
- public Response getResponse(InternalRequest internalRequest) {
+ public CompletableFuture getResponse(InternalRequest internalRequest) {
@Untrusted(reason = "from header") String accessAddress = internalRequest.getAccessAddress(webserverConfiguration);
@Untrusted String requestedPath = internalRequest.getRequestedPathAndQuery();
boolean blocked = false;
- Response response;
+ CompletableFuture response;
@Untrusted Request request = null;
if (bruteForceGuard.shouldPreventRequest(accessAddress)) {
- response = responseFactory.failedLoginAttempts403();
+ response = CompletableFuture.completedFuture(responseFactory.failedLoginAttempts403());
blocked = true;
} else if (rateLimitGuard.shouldPreventRequest(requestedPath, accessAddress)) {
- response = responseFactory.failedRateLimit403();
+ response = CompletableFuture.completedFuture(responseFactory.failedRateLimit403());
blocked = true;
} else if (!webserverConfiguration.getAllowedIpList().isAllowed(accessAddress)) {
webserverConfiguration.getWebserverLogMessages()
.warnAboutWhitelistBlock(accessAddress, internalRequest.getRequestedURIString());
- response = responseFactory.ipWhitelist403(accessAddress);
+ response = CompletableFuture.completedFuture(responseFactory.ipWhitelist403(accessAddress));
} else {
- try {
- request = internalRequest.toRequest(accessAddress);
- response = attemptToResolve(request, accessAddress);
- } catch (WebUserAuthException thrownByAuthentication) {
- response = processFailedAuthentication(internalRequest, accessAddress, thrownByAuthentication);
- }
+ request = internalRequest.toRequest(accessAddress);
+ response = attemptToResolve(request, accessAddress)
+ .exceptionallyCompose(throwable -> {
+ Throwable cause = throwable;
+ while (cause instanceof java.util.concurrent.CompletionException || cause instanceof java.util.concurrent.ExecutionException) {
+ if (cause.getCause() != null) {
+ cause = cause.getCause();
+ } else {
+ break;
+ }
+ }
+ if (cause instanceof WebUserAuthException thrownByAuthentication) {
+ return CompletableFuture.completedFuture(
+ processFailedAuthentication(internalRequest, accessAddress, thrownByAuthentication));
+ }
+ return CompletableFuture.failedFuture(throwable);
+ });
}
- response.getHeaders().putIfAbsent("Access-Control-Allow-Origin", webserverConfiguration.getAllowedCorsOrigin());
- response.getHeaders().putIfAbsent("Access-Control-Allow-Methods", "GET, OPTIONS");
- response.getHeaders().putIfAbsent("Access-Control-Allow-Credentials", "true");
- response.getHeaders().putIfAbsent("X-Robots-Tag", "noindex, nofollow");
+ response.whenComplete((r, e) -> {
+ if (r != null) {
+ Map headers = r.getHeaders();
+ headers.putIfAbsent("Access-Control-Allow-Origin", webserverConfiguration.getAllowedCorsOrigin());
+ headers.putIfAbsent("Access-Control-Allow-Methods", "GET, OPTIONS");
+ headers.putIfAbsent("Access-Control-Allow-Credentials", "true");
+ headers.putIfAbsent("X-Robots-Tag", "noindex, nofollow");
+ }
+ });
if (!blocked) {
- accessLogger.log(internalRequest, request, response);
+ final Request doneRequest = request;
+ response.whenComplete((r, e) -> accessLogger.log(internalRequest, doneRequest, r));
}
return response;
}
- private Response attemptToResolve(@Untrusted Request request, @Untrusted String accessAddress) {
- Response response = protocolUpgradeResponse(request)
- .orElseGet(() -> responseResolver.getResponse(request));
- request.getUser().ifPresent(user -> processSuccessfulLogin(response.getCode(), accessAddress));
+ private CompletableFuture attemptToResolve(@Untrusted Request request, @Untrusted String accessAddress) {
+ CompletableFuture response;
+ try {
+ response = protocolUpgradeResponse(request)
+ .orElseGet(() -> responseResolver.getResponse(request));
+ } catch (Throwable t) {
+ response = responseResolver.handleException(request, t);
+ }
+ response.whenComplete((r, e) -> {
+ if (r != null) {
+ request.getUser().ifPresent(user -> processSuccessfulLogin(r.getCode(), accessAddress));
+ }
+ });
return response;
}
- private Optional protocolUpgradeResponse(@Untrusted Request request) {
+ private Optional> protocolUpgradeResponse(@Untrusted Request request) {
@Untrusted Optional upgrade = request.getHeader(HttpHeader.UPGRADE.asString());
if (upgrade.isPresent()) {
@Untrusted String value = upgrade.get();
if ("h2c".equals(value) || "h2".equals(value)) {
- return Optional.of(Response.builder()
+ return Optional.of(CompletableFuture.completedFuture(Response.builder()
.setStatus(101)
.setHeader("Connection", HttpHeader.UPGRADE.asString())
.setHeader(HttpHeader.UPGRADE.asString(), value)
- .build());
+ .build()));
}
}
return Optional.empty();
diff --git a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/resolver/json/query/QueryJSONResolver.java b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/resolver/json/query/QueryJSONResolver.java
index 7d366b2f6c..6d5a025e4d 100644
--- a/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/resolver/json/query/QueryJSONResolver.java
+++ b/Plan/common/src/main/java/com/djrapitops/plan/delivery/webserver/resolver/json/query/QueryJSONResolver.java
@@ -26,8 +26,8 @@
import com.djrapitops.plan.delivery.formatting.Formatters;
import com.djrapitops.plan.delivery.rendering.json.PlayersTableJSONCreator;
import com.djrapitops.plan.delivery.rendering.json.graphs.GraphJSONCreator;
+import com.djrapitops.plan.delivery.web.resolver.AsyncResolver;
import com.djrapitops.plan.delivery.web.resolver.MimeType;
-import com.djrapitops.plan.delivery.web.resolver.Resolver;
import com.djrapitops.plan.delivery.web.resolver.Response;
import com.djrapitops.plan.delivery.web.resolver.exception.BadRequestException;
import com.djrapitops.plan.delivery.web.resolver.request.Request;
@@ -37,6 +37,7 @@
import com.djrapitops.plan.extension.implementation.storage.queries.ExtensionQueryResultTableDataQuery;
import com.djrapitops.plan.identification.ServerInfo;
import com.djrapitops.plan.identification.ServerUUID;
+import com.djrapitops.plan.processing.Processing;
import com.djrapitops.plan.settings.config.PlanConfig;
import com.djrapitops.plan.settings.config.paths.DisplaySettings;
import com.djrapitops.plan.settings.config.paths.TimeSettings;
@@ -70,13 +71,15 @@
import java.nio.charset.StandardCharsets;
import java.text.ParseException;
import java.util.*;
+import java.util.concurrent.CompletableFuture;
@Singleton
@Path("/v1/query")
-public class QueryJSONResolver implements Resolver {
+public class QueryJSONResolver implements AsyncResolver {
private final QueryFilters filters;
+ private final Processing processing;
private final PlanConfig config;
private final DBSystem dbSystem;
private final ServerInfo serverInfo;
@@ -89,6 +92,7 @@ public class QueryJSONResolver implements Resolver {
@Inject
public QueryJSONResolver(
QueryFilters filters,
+ Processing processing,
PlanConfig config,
DBSystem dbSystem,
ServerInfo serverInfo, JSONStorage jsonStorage,
@@ -98,6 +102,7 @@ public QueryJSONResolver(
Gson gson
) {
this.filters = filters;
+ this.processing = processing;
this.config = config;
this.dbSystem = dbSystem;
this.serverInfo = serverInfo;
@@ -134,32 +139,38 @@ public boolean canAccess(Request request) {
requestBody = @RequestBody(content = @Content(schema = @Schema(implementation = InputQueryDto.class)))
)
@Override
- public Optional resolve(Request request) {
- return Optional.of(getResponse(request));
+ public CompletableFuture> resolveAsync(Request request) {
+ return getResponse(request);
}
- private Response getResponse(@Untrusted Request request) {
+ private CompletableFuture> getResponse(@Untrusted Request request) {
Optional user = request.getUser();
boolean canAccessCache = user.map(u -> u.hasPermission(WebPermission.ACCESS_QUERY)).orElse(true);
- Optional cachedResult = canAccessCache ? checkForCachedResult(request) : Optional.empty();
- if (cachedResult.isPresent()) return cachedResult.get();
-
- InputQueryDto inputQuery = parseInputQuery(request);
- @Untrusted List queries = inputQuery.getFilters();
-
- // Check user has permission for the filter if login is enabled.
- if (user.isPresent()) {
- Optional errorResponse = checkFilterPermissions(queries, user.get());
- if (errorResponse.isPresent()) {
- return errorResponse.get();
+ if (canAccessCache) {
+ Optional cached = checkForCachedResult(request);
+ if (cached.isPresent()) {
+ return CompletableFuture.completedFuture(cached);
}
}
- Filter.Result result = filters.apply(queries);
- List resultPath = result.getInverseResultPath();
- Collections.reverse(resultPath);
+ return request.getRequestBodyAsync().thenCompose(bodyBytes -> {
+ InputQueryDto inputQuery = parseInputQuery(request, bodyBytes);
+ List queries = inputQuery.getFilters();
+
+ if (user.isPresent()) {
+ Optional errorResponse = checkFilterPermissions(queries, user.get());
+ if (errorResponse.isPresent()) {
+ return CompletableFuture.completedFuture(errorResponse);
+ }
+ }
- return buildAndStoreResponse(inputQuery, result, resultPath);
+ return CompletableFuture.supplyAsync(() -> {
+ Filter.Result result = filters.apply(queries);
+ List resultPath = result.getInverseResultPath();
+ Collections.reverse(resultPath);
+ return Optional.of(buildAndStoreResponse(inputQuery, result, resultPath));
+ }, processing.getNonCriticalExecutor());
+ });
}
private Optional checkFilterPermissions(List queries, WebUser user) {
@@ -203,11 +214,11 @@ private WebPermission[] getAllowingPermissions(@Untrusted String filterKind) {
}
}
- private InputQueryDto parseInputQuery(@Untrusted Request request) {
- if (request.getRequestBody().length == 0) {
+ private InputQueryDto parseInputQuery(@Untrusted Request request, byte[] bodyBytes) {
+ if (bodyBytes.length == 0) {
return parseInputQueryFromQueryParams(request);
} else {
- return RequestBodyConverter.bodyJson(request, gson, InputQueryDto.class);
+ return RequestBodyConverter.bodyJson(bodyBytes, gson, InputQueryDto.class);
}
}
diff --git a/Plan/common/src/main/java/com/djrapitops/plan/processing/Processing.java b/Plan/common/src/main/java/com/djrapitops/plan/processing/Processing.java
index 96d07a4ec6..e889f6e745 100644
--- a/Plan/common/src/main/java/com/djrapitops/plan/processing/Processing.java
+++ b/Plan/common/src/main/java/com/djrapitops/plan/processing/Processing.java
@@ -213,4 +213,8 @@ private void ensureShutdown() {
public Executor getCriticalExecutor() {
return criticalExecutor;
}
+
+ public ExecutorService getNonCriticalExecutor() {
+ return nonCriticalExecutor;
+ }
}
diff --git a/Plan/common/src/main/java/com/djrapitops/plan/storage/database/Database.java b/Plan/common/src/main/java/com/djrapitops/plan/storage/database/Database.java
index af91b41608..92762eecae 100644
--- a/Plan/common/src/main/java/com/djrapitops/plan/storage/database/Database.java
+++ b/Plan/common/src/main/java/com/djrapitops/plan/storage/database/Database.java
@@ -28,6 +28,7 @@
import java.sql.SQLException;
import java.util.*;
import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutorService;
import java.util.function.Supplier;
/**
@@ -163,6 +164,8 @@ default Sql getSql() {
int getTransactionQueueSize();
+ Async async();
+
/**
* Possible State changes:
* CLOSED to PATCHING (Database init),
@@ -177,4 +180,32 @@ enum State {
OPEN,
CLOSING
}
+
+ public static class Async {
+ private final Database database;
+ private final ExecutorService executorService;
+
+ public Async(Database database, ExecutorService executorService) {
+ this.database = database;
+ this.executorService = executorService;
+ }
+
+ public CompletableFuture query(Query query) {
+ return CompletableFuture.supplyAsync(() -> database.query(query), executorService);
+ }
+
+ public CompletableFuture> queryOptional(String sql, RowExtractor rowExtractor, Object... parameters) {
+ return CompletableFuture.supplyAsync(() -> database.queryOptional(sql, rowExtractor, parameters), executorService);
+ }
+
+ public CompletableFuture> queryList(String sql, RowExtractor rowExtractor, Object... parameters) {return CompletableFuture.supplyAsync(() -> database.queryList(sql, rowExtractor, parameters), executorService);}
+
+ public CompletableFuture> querySet(String sql, RowExtractor rowExtractor, Object... parameters) {return CompletableFuture.supplyAsync(() -> database.querySet(sql, rowExtractor, parameters), executorService);}
+
+ public , T> CompletableFuture queryCollection(String sql, RowExtractor rowExtractor, Supplier collectionConstructor, Object... parameters) {return CompletableFuture.supplyAsync(() -> database.queryCollection(sql, rowExtractor, collectionConstructor, parameters), executorService);}
+
+ public CompletableFuture