Make request-path reactive-safe: audit/job aspects, JobExecutorService, license singleton race
This commit is contained in:
@@ -9,12 +9,13 @@ import java.util.function.Supplier;
|
||||
|
||||
import org.slf4j.MDC;
|
||||
|
||||
import io.quarkus.vertx.http.runtime.CurrentVertxRequest;
|
||||
|
||||
import jakarta.annotation.Priority;
|
||||
import jakarta.inject.Inject;
|
||||
import jakarta.interceptor.AroundInvoke;
|
||||
import jakarta.interceptor.Interceptor;
|
||||
import jakarta.interceptor.InvocationContext;
|
||||
import jakarta.servlet.http.HttpServletRequest;
|
||||
|
||||
import lombok.extern.slf4j.Slf4j;
|
||||
|
||||
@@ -41,29 +42,66 @@ public class AutoJobAspect {
|
||||
private static final Duration RETRY_BASE_DELAY = Duration.ofMillis(100);
|
||||
|
||||
private final JobExecutorService jobExecutorService;
|
||||
private final HttpServletRequest request;
|
||||
// Reactive-safe access to the current request. The undertow HttpServletRequest proxy throws
|
||||
// UT000048 ("No request is currently active") on RESTEasy Reactive worker threads, so query
|
||||
// params / method / path / attributes are read from the Vert.x request instead, degrading to
|
||||
// null/empty when no request is active.
|
||||
private final CurrentVertxRequest currentVertxRequest;
|
||||
private final FileStorage fileStorage;
|
||||
|
||||
@Inject
|
||||
public AutoJobAspect(
|
||||
JobExecutorService jobExecutorService,
|
||||
HttpServletRequest request,
|
||||
CurrentVertxRequest currentVertxRequest,
|
||||
FileStorage fileStorage) {
|
||||
this.jobExecutorService = jobExecutorService;
|
||||
this.request = request;
|
||||
this.currentVertxRequest = currentVertxRequest;
|
||||
this.fileStorage = fileStorage;
|
||||
}
|
||||
|
||||
private io.vertx.core.http.HttpServerRequest vertxRequest() {
|
||||
try {
|
||||
var current = currentVertxRequest.getCurrent();
|
||||
return current != null ? current.request() : null;
|
||||
} catch (RuntimeException e) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
private String requestParam(String name) {
|
||||
io.vertx.core.http.HttpServerRequest req = vertxRequest();
|
||||
return req != null ? req.getParam(name) : null;
|
||||
}
|
||||
|
||||
private String requestMethod() {
|
||||
io.vertx.core.http.HttpServerRequest req = vertxRequest();
|
||||
return req != null ? req.method().name() : "";
|
||||
}
|
||||
|
||||
private String requestUri() {
|
||||
io.vertx.core.http.HttpServerRequest req = vertxRequest();
|
||||
return req != null ? req.path() : "";
|
||||
}
|
||||
|
||||
private Object requestAttribute(String name) {
|
||||
try {
|
||||
var current = currentVertxRequest.getCurrent();
|
||||
return current != null ? current.get(name) : null;
|
||||
} catch (RuntimeException e) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
@AroundInvoke
|
||||
public Object wrapWithJobExecution(InvocationContext ctx) throws Exception {
|
||||
AutoJobPostMapping autoJobPostMapping =
|
||||
ctx.getMethod().getAnnotation(AutoJobPostMapping.class);
|
||||
// Extract parameters from the request and annotation
|
||||
boolean async = Boolean.parseBoolean(request.getParameter("async"));
|
||||
boolean async = Boolean.parseBoolean(requestParam("async"));
|
||||
log.debug(
|
||||
"AutoJobAspect: Processing {} {} with async={}",
|
||||
request.getMethod(),
|
||||
request.getRequestURI(),
|
||||
requestMethod(),
|
||||
requestUri(),
|
||||
async);
|
||||
long timeout = autoJobPostMapping.timeout();
|
||||
int retryCount = autoJobPostMapping.retryCount();
|
||||
@@ -310,7 +348,7 @@ public class AutoJobAspect {
|
||||
|
||||
private String getJobIdFromContext() {
|
||||
try {
|
||||
return (String) request.getAttribute("jobId");
|
||||
return (String) requestAttribute("jobId");
|
||||
} catch (Exception e) {
|
||||
log.debug("Could not retrieve job ID from context: {}", e.getMessage());
|
||||
return null;
|
||||
|
||||
@@ -13,7 +13,6 @@ import org.eclipse.microprofile.config.inject.ConfigProperty;
|
||||
|
||||
import jakarta.enterprise.context.ApplicationScoped;
|
||||
import jakarta.enterprise.inject.Instance;
|
||||
import jakarta.servlet.http.HttpServletRequest;
|
||||
import jakarta.ws.rs.core.HttpHeaders;
|
||||
import jakarta.ws.rs.core.MediaType;
|
||||
import jakarta.ws.rs.core.Response;
|
||||
@@ -36,7 +35,10 @@ public class JobExecutorService {
|
||||
|
||||
private final TaskManager taskManager;
|
||||
private final FileStorage fileStorage;
|
||||
private final HttpServletRequest request;
|
||||
// Reactive-safe: the undertow HttpServletRequest proxy throws UT000048 on RESTEasy Reactive
|
||||
// worker threads, so the per-request "jobId" attribute is stored on the Vert.x RoutingContext
|
||||
// instead (read back via AutoJobAspect). Off a live request this degrades to a no-op.
|
||||
private final io.quarkus.vertx.http.runtime.CurrentVertxRequest currentVertxRequest;
|
||||
private final ResourceMonitor resourceMonitor;
|
||||
private final JobQueue jobQueue;
|
||||
private final ExecutorService executor = ExecutorFactory.newVirtualThreadExecutor();
|
||||
@@ -47,7 +49,7 @@ public class JobExecutorService {
|
||||
public JobExecutorService(
|
||||
TaskManager taskManager,
|
||||
FileStorage fileStorage,
|
||||
HttpServletRequest request,
|
||||
io.quarkus.vertx.http.runtime.CurrentVertxRequest currentVertxRequest,
|
||||
ResourceMonitor resourceMonitor,
|
||||
JobQueue jobQueue,
|
||||
@ConfigProperty(name = "spring.mvc.async.request-timeout", defaultValue = "1200000")
|
||||
@@ -56,7 +58,7 @@ public class JobExecutorService {
|
||||
String sessionTimeout) {
|
||||
this.taskManager = taskManager;
|
||||
this.fileStorage = fileStorage;
|
||||
this.request = request;
|
||||
this.currentVertxRequest = currentVertxRequest;
|
||||
this.resourceMonitor = resourceMonitor;
|
||||
this.jobQueue = jobQueue;
|
||||
|
||||
@@ -85,8 +87,13 @@ public class JobExecutorService {
|
||||
|
||||
log.debug("Generated jobId: {} (base: {})", scopedJobKey, baseJobId);
|
||||
|
||||
if (request != null) {
|
||||
request.setAttribute("jobId", scopedJobKey);
|
||||
try {
|
||||
var current = currentVertxRequest.getCurrent();
|
||||
if (current != null) {
|
||||
current.put("jobId", scopedJobKey);
|
||||
}
|
||||
} catch (RuntimeException ignored) {
|
||||
// No active request (e.g. async/background execution) - jobId attribute is optional.
|
||||
}
|
||||
|
||||
String jobId = scopedJobKey;
|
||||
|
||||
@@ -98,7 +98,10 @@ public class AuditAspect {
|
||||
// an
|
||||
// HTTP request scope the injected proxy resolves to null, so we treat a null request the
|
||||
// same way the original treated a null ServletRequestAttributes.
|
||||
HttpServletRequest req = request;
|
||||
// Reactive-safe: the injected proxy is non-null but throws UT000048 when touched off an
|
||||
// active servlet request (RESTEasy Reactive worker threads). Resolve via the guarded
|
||||
// AuditService.getCurrentRequest(), which returns null outside a live servlet request.
|
||||
HttpServletRequest req = auditService.getCurrentRequest();
|
||||
boolean isHttpRequest = req != null;
|
||||
|
||||
String capturedIp = MDC.get("auditIp");
|
||||
|
||||
+28
-6
@@ -95,10 +95,29 @@ public class ControllerAuditAspect {
|
||||
*/
|
||||
@AroundInvoke
|
||||
public Object auditEndpoint(InvocationContext ctx) throws Throwable {
|
||||
String httpMethod = request != null ? request.getMethod() : "POST";
|
||||
// Reactive-safe: the injected HttpServletRequest proxy is never null but throws UT000048
|
||||
// ("No request is currently active") when touched on a RESTEasy Reactive worker thread.
|
||||
// Resolve the verb through the guarded AuditService.getCurrentRequest() (returns null off a
|
||||
// servlet request) and fall back to POST, mirroring the original non-web behaviour.
|
||||
HttpServletRequest current = auditService.getCurrentRequest();
|
||||
String httpMethod = current != null ? current.getMethod() : "POST";
|
||||
return auditController(ctx, httpMethod != null ? httpMethod : "POST");
|
||||
}
|
||||
|
||||
// Reactive-safe accessor for the response proxy: touching it off an active servlet request
|
||||
// throws UT000048, so treat that (and an unsatisfied proxy) as "no response available".
|
||||
private HttpServletResponse safeResponse() {
|
||||
try {
|
||||
if (response == null) {
|
||||
return null;
|
||||
}
|
||||
response.getStatus();
|
||||
return response;
|
||||
} catch (RuntimeException e) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
private Object auditController(InvocationContext joinPoint, String httpMethod)
|
||||
throws Throwable {
|
||||
Method method = joinPoint.getMethod();
|
||||
@@ -135,8 +154,8 @@ public class ControllerAuditAspect {
|
||||
}
|
||||
}
|
||||
|
||||
HttpServletRequest req = request;
|
||||
HttpServletResponse resp = response;
|
||||
HttpServletRequest req = auditService.getCurrentRequest();
|
||||
HttpServletResponse resp = safeResponse();
|
||||
|
||||
String previousPrincipal = MDC.get("auditPrincipal");
|
||||
String previousOrigin = MDC.get("auditOrigin");
|
||||
@@ -271,9 +290,12 @@ public class ControllerAuditAspect {
|
||||
// Using AuditUtils.determineAuditEventType instead
|
||||
|
||||
private String getRequestPath(Method method, String httpMethod) {
|
||||
// Prefer actual request URI over annotation patterns (which may contain regex)
|
||||
if (request != null) {
|
||||
return request.getRequestURI();
|
||||
// Prefer actual request URI over annotation patterns (which may contain regex).
|
||||
// Reactive-safe: go through the guarded accessor (the raw proxy throws UT000048
|
||||
// off-thread).
|
||||
HttpServletRequest current = auditService.getCurrentRequest();
|
||||
if (current != null) {
|
||||
return current.getRequestURI();
|
||||
}
|
||||
// Fallback: try JAX-RS @Path annotation on method/class; return empty string if not present
|
||||
// TODO: Migration required - resolve path from jakarta.ws.rs.@Path on the declaring class
|
||||
|
||||
+41
-13
@@ -10,6 +10,8 @@ import java.util.UUID;
|
||||
import javax.crypto.Mac;
|
||||
import javax.crypto.spec.SecretKeySpec;
|
||||
|
||||
import io.quarkus.narayana.jta.QuarkusTransaction;
|
||||
|
||||
import jakarta.enterprise.context.ApplicationScoped;
|
||||
import jakarta.enterprise.inject.Instance;
|
||||
import jakarta.transaction.Transactional;
|
||||
@@ -56,23 +58,49 @@ public class UserLicenseSettingsService {
|
||||
*
|
||||
* @return The current settings
|
||||
*/
|
||||
// Serializes singleton creation within this JVM so two concurrent callers (e.g. the startup
|
||||
// license sync racing with the first inbound request) cannot both INSERT id=1.
|
||||
private static final Object CREATE_LOCK = new Object();
|
||||
|
||||
@Transactional
|
||||
public UserLicenseSettings getOrCreateSettings() {
|
||||
Optional<UserLicenseSettings> existing = settingsRepository.findSettings();
|
||||
if (existing.isPresent()) {
|
||||
return existing.get();
|
||||
}
|
||||
// MIGRATION: Spring Data save() on this manually-@Id'd singleton did an upsert; the
|
||||
// migrated
|
||||
// persist() is INSERT-only and a PK violation when two transactions create id=1 at once
|
||||
// (startup sync vs first request). Create the row in its OWN committed transaction, guarded
|
||||
// by a JVM lock with a fresh-tx re-check, so exactly one INSERT happens; then reload it
|
||||
// into
|
||||
// the current transaction so callers that modify+persist operate on a managed entity.
|
||||
synchronized (CREATE_LOCK) {
|
||||
boolean alreadyCreated =
|
||||
QuarkusTransaction.requiringNew()
|
||||
.call(() -> settingsRepository.findSettings().isPresent());
|
||||
if (!alreadyCreated) {
|
||||
log.info("Initializing user license settings");
|
||||
QuarkusTransaction.requiringNew()
|
||||
.run(
|
||||
() -> {
|
||||
UserLicenseSettings settings = new UserLicenseSettings();
|
||||
settings.setId(UserLicenseSettings.SINGLETON_ID);
|
||||
settings.setGrandfatheredUserCount(0);
|
||||
settings.setLicenseMaxUsers(0);
|
||||
settings.setGrandfatheringLocked(false);
|
||||
settings.setIntegritySalt(UUID.randomUUID().toString());
|
||||
settings.setGrandfatheredUserSignature("");
|
||||
settingsRepository.persist(settings);
|
||||
});
|
||||
}
|
||||
}
|
||||
return settingsRepository
|
||||
.findSettings()
|
||||
.orElseGet(
|
||||
() -> {
|
||||
log.info("Initializing user license settings");
|
||||
UserLicenseSettings settings = new UserLicenseSettings();
|
||||
settings.setId(UserLicenseSettings.SINGLETON_ID);
|
||||
settings.setGrandfatheredUserCount(0);
|
||||
settings.setLicenseMaxUsers(0);
|
||||
settings.setGrandfatheringLocked(false);
|
||||
settings.setIntegritySalt(UUID.randomUUID().toString());
|
||||
settings.setGrandfatheredUserSignature("");
|
||||
settingsRepository.persist(settings);
|
||||
return settings;
|
||||
});
|
||||
.orElseThrow(
|
||||
() ->
|
||||
new IllegalStateException(
|
||||
"User license settings missing immediately after creation"));
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
Reference in New Issue
Block a user