From 6120b64ff2bfe31d3d2072c4c243e051b984fc61 Mon Sep 17 00:00:00 2001 From: Stephen Lloyd Date: Fri, 26 Jun 2026 10:38:05 +0100 Subject: [PATCH 1/6] Creates a transactional asserted object to allow the jobs to be created Added the persist package of UWS to the hibernate objects in application.properties for UWSJobEntity visibility --- .../org/javastro/ivoa/tap/QueryResource.java | 16 +++++----- .../org/javastro/ivoa/tap/TAPJobService.java | 22 ++++++++++++++ .../javastro/ivoa/tap/TapConfiguration.java | 29 +++++++++++++++++-- src/main/resources/application.properties | 4 +-- 4 files changed, 58 insertions(+), 13 deletions(-) create mode 100644 src/main/java/org/javastro/ivoa/tap/TAPJobService.java diff --git a/src/main/java/org/javastro/ivoa/tap/QueryResource.java b/src/main/java/org/javastro/ivoa/tap/QueryResource.java index dace59a..50cf491 100644 --- a/src/main/java/org/javastro/ivoa/tap/QueryResource.java +++ b/src/main/java/org/javastro/ivoa/tap/QueryResource.java @@ -10,6 +10,7 @@ import io.smallrye.mutiny.infrastructure.Infrastructure; import jakarta.enterprise.context.ApplicationScoped; import jakarta.inject.Inject; +import jakarta.transaction.Transactional; import jakarta.ws.rs.*; import jakarta.ws.rs.core.Context; import jakarta.ws.rs.core.MediaType; @@ -32,11 +33,7 @@ import uk.ac.starlink.table.*; -import java.io.IOException; -import java.net.URI; -import java.nio.file.Files; import java.time.Duration; -import java.util.Map; /** * Main TAP Query. @@ -56,6 +53,9 @@ public class QueryResource { @Inject TAPHelper tapHelper; + @Inject + TAPJobService jobService; + private static final Logger log = LoggerFactory.getLogger(QueryResource.class); @GET @@ -82,7 +82,6 @@ public Uni syncPost(@RestForm("QUERY") String query, @RestFo return handleJob(query, lang, responseformat, maxrec, runid, upload, input, uriInfo); } - private Uni handleJob(String query, String lang, String responseformat, Long maxrec, String runid, String upload, MultipartFormDataInput input, UriInfo uriInfo) { final Duration SYNC_WAIT = Duration.ofSeconds(syncTimeoutSeconds); return Uni.createFrom().deferred(() -> { @@ -92,9 +91,10 @@ private Uni handleJob(String query, String lang, String resp if(upload != null && !upload.isEmpty() ) { tapUploader = new QuarkusTapUploader(upload, input); } - job = (TAPJob) tapHelper.jobmanager.createJob( - new TAPJobSpecification(query, lang, responseformat, maxrec, runid, tapUploader) - ); + // job = (TAPJob) tapHelper.jobmanager.createJob( + // new TAPJobSpecification(query, lang, responseformat, maxrec, runid, tapUploader) + // ); + job = jobService.createJob(new TAPJobSpecification(query, lang, responseformat, maxrec, runid, tapUploader)); tapHelper.jobmanager.runJob(job.getID()); // automatically run the job } catch (UWSException e) { diff --git a/src/main/java/org/javastro/ivoa/tap/TAPJobService.java b/src/main/java/org/javastro/ivoa/tap/TAPJobService.java new file mode 100644 index 0000000..d72c6dd --- /dev/null +++ b/src/main/java/org/javastro/ivoa/tap/TAPJobService.java @@ -0,0 +1,22 @@ +package org.javastro.ivoa.tap; + +import jakarta.enterprise.context.ApplicationScoped; +import jakarta.transaction.Transactional; +import org.javastro.ivoacore.tap.TAPJob; +import org.javastro.ivoacore.tap.TAPJobSpecification; +import org.javastro.ivoacore.uws.UWSException; + +@ApplicationScoped +public class TAPJobService { + + private final TAPHelper tapHelper; + + public TAPJobService(TAPHelper tapHelper) { + this.tapHelper = tapHelper; + } + + @Transactional + public TAPJob createJob(TAPJobSpecification spec) throws UWSException { + return (TAPJob) tapHelper.jobmanager.createJob(spec); + } +} diff --git a/src/main/java/org/javastro/ivoa/tap/TapConfiguration.java b/src/main/java/org/javastro/ivoa/tap/TapConfiguration.java index 238a27b..84f63e8 100644 --- a/src/main/java/org/javastro/ivoa/tap/TapConfiguration.java +++ b/src/main/java/org/javastro/ivoa/tap/TapConfiguration.java @@ -10,18 +10,23 @@ import jakarta.enterprise.inject.Produces; import jakarta.inject.Inject; import jakarta.inject.Singleton; +import jakarta.persistence.EntityManager; import org.eclipse.microprofile.config.inject.ConfigProperty; import org.javastro.ivoa.entities.resource.Capability; import org.javastro.ivoa.entities.vosi.capabilities.Capabilities; import org.javastro.ivoacore.common.ServiceLocator; import org.javastro.ivoacore.tap.TAPJob; +import org.javastro.ivoacore.tap.TAPJobSpecification; import org.javastro.ivoacore.tap.schema.SchemaProvider; import org.javastro.ivoacore.tap.schema.VODMLSchemaProvider; +import org.javastro.ivoacore.uws.JobFactoryAggregator; import org.javastro.ivoacore.uws.JobManager; import org.javastro.ivoacore.uws.environment.DefaultEnvironmentFactory; import org.javastro.ivoacore.uws.environment.DefaultExecutionPolicy; import org.javastro.ivoacore.uws.environment.EnvironmentFactory; -import org.javastro.ivoacore.uws.persist.MemoryBasedJobStore; +import org.javastro.ivoacore.uws.persist.CachedJobStore; +import org.javastro.ivoacore.uws.persist.DatabaseJobStore; +import org.javastro.ivoacore.uws.persist.JobStore; import org.javastro.ivoacore.vosi.CapabilityBuilder; import org.javastro.ivoacore.vosi.VOSIProvider; @@ -46,6 +51,9 @@ public class TapConfiguration { @Inject DataSource ds; + @Inject + EntityManager em; + @ConfigProperty(name="ivoa.tap.dbCaseSensitive", defaultValue = "false") boolean isDbCaseSensitive; @@ -103,8 +111,23 @@ JobManager uws(SchemaProvider schemaProvider) { } EnvironmentFactory env = new DefaultEnvironmentFactory(tmpdir); - MemoryBasedJobStore store = new MemoryBasedJobStore(); + // MemoryBasedJobStore store = new MemoryBasedJobStore(); + + TAPJob.JobFactory tapJobFactory = new TAPJob.JobFactory(ds, schemaProvider, env); + + JobFactoryAggregator factoryAggregator = new JobFactoryAggregator(); + factoryAggregator.addFactory(tapJobFactory); + + JobStore store = new CachedJobStore( + DatabaseJobStore.forJobType( + em, + TAPJobSpecification.class, + "TAP", + factoryAggregator + ) + ); + //CachedJobStore store = new CachedJobStore(new DatabaseJobStore(ds, em, )); DefaultExecutionPolicy policy = new DefaultExecutionPolicy(); - return new JobManager(new TAPJob.JobFactory(ds, schemaProvider, env), store, policy); + return new JobManager(factoryAggregator, store, policy); } } diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties index 91330d2..7e02119 100644 --- a/src/main/resources/application.properties +++ b/src/main/resources/application.properties @@ -30,7 +30,7 @@ quarkus.hibernate-orm.dialect=org.javastro.ivoacore.pgsphere.PgSphereDialect #ORM setup - mainly for the TAPSchema -quarkus.hibernate-orm.packages=org.ivoa.dm,org.ivoa.vodml.stdtypes,org.javastro.ivoacore.pgsphere,org.javastro.ivoa.tap.entities +quarkus.hibernate-orm.packages=org.ivoa.dm,org.ivoa.vodml.stdtypes,org.javastro.ivoacore.pgsphere,org.javastro.ivoa.tap.entities,org.javastro.ivoacore.uws.persist #%dev,%test.quarkus.hibernate-orm.database.generation=drop-and-create %prod.quarkus.hibernate-orm.database.generation=update quarkus.hibernate-orm.database.generation.create-schemas=true @@ -69,4 +69,4 @@ quarkus.kubernetes.ingress.path-type=Prefix #attempt to reduce logging noise from netty websocket -quarkus.log.category."io.netty.handler.codec.http.websocketx".use-parent-handlers=false \ No newline at end of file +quarkus.log.category."io.netty.handler.codec.http.websocketx".use-parent-handlers=false From 523138b09e99bb8df77f17d480ee908fa954f3f1 Mon Sep 17 00:00:00 2001 From: Stephen Lloyd Date: Fri, 26 Jun 2026 16:15:30 +0100 Subject: [PATCH 2/6] Had to add overrides for async tasks as they need to be transactional Same for the service object used on the sync methods as they exist in their own right (due to the Uni launching the create job on a separate thread) --- .../javastro/ivoa/tap/AsyncQueryResource.java | 23 ++++++++++++++++--- .../org/javastro/ivoa/tap/QueryResource.java | 1 - .../org/javastro/ivoa/tap/TAPJobService.java | 5 ++++ 3 files changed, 25 insertions(+), 4 deletions(-) diff --git a/src/main/java/org/javastro/ivoa/tap/AsyncQueryResource.java b/src/main/java/org/javastro/ivoa/tap/AsyncQueryResource.java index 197bf5e..80b47a8 100644 --- a/src/main/java/org/javastro/ivoa/tap/AsyncQueryResource.java +++ b/src/main/java/org/javastro/ivoa/tap/AsyncQueryResource.java @@ -8,6 +8,7 @@ import jakarta.enterprise.context.ApplicationScoped; import jakarta.inject.Inject; +import jakarta.transaction.Transactional; import jakarta.ws.rs.*; import jakarta.ws.rs.core.*; import org.eclipse.microprofile.openapi.annotations.tags.Tag; @@ -23,9 +24,6 @@ import org.jboss.resteasy.reactive.RestResponse; import org.jboss.resteasy.reactive.server.multipart.MultipartFormDataInput; -import java.net.URI; -import java.util.Map; - /** * Main Async TAP Query. * Created on 04/03/2026 by Paul Harrison (paul.harrison@manchester.ac.uk). @@ -58,6 +56,7 @@ protected Response redirectToJob(String jobid) { //IMPL the two query endpoints are in different resources for routing purposes @POST @Consumes({MediaType.APPLICATION_FORM_URLENCODED, MediaType.MULTIPART_FORM_DATA}) + @Transactional public Response async(@RestForm("QUERY") String query, @RestForm("LANG") String lang, @RestForm("RESPONSEFORMAT") String responseformat, @RestForm("MAXREC") Long maxrec, @RestForm("RUNID") String runid, @RestForm("UPLOAD") String upload, MultipartFormDataInput input, @Context UriInfo uriInfo) throws UWSException { @@ -78,4 +77,22 @@ public RestResponse getVotable(@PathParam("jobid") String jo .header(HttpHeaders.CONTENT_DISPOSITION, "result.vot") .build(); } + + //----------------------- Need to make database modifying operations transactional ---------------------------------- + // Which means the base UWS modifying tasks need to be wrapped in a transactional override + @Override + @DELETE + @Path("{jobid}") + @Transactional + public Response deleteJob(@PathParam("jobid")String jobid) throws UWSException { + return super.deleteJob(jobid); + } + + @Override + @POST + @Path("{jobid}/phase") + @Transactional + public Response setPhase(@PathParam("jobid") String jobid, @FormParam("PHASE") String phase) throws UWSException { + return super.setPhase(jobid, phase); + } } diff --git a/src/main/java/org/javastro/ivoa/tap/QueryResource.java b/src/main/java/org/javastro/ivoa/tap/QueryResource.java index 50cf491..f84cce5 100644 --- a/src/main/java/org/javastro/ivoa/tap/QueryResource.java +++ b/src/main/java/org/javastro/ivoa/tap/QueryResource.java @@ -10,7 +10,6 @@ import io.smallrye.mutiny.infrastructure.Infrastructure; import jakarta.enterprise.context.ApplicationScoped; import jakarta.inject.Inject; -import jakarta.transaction.Transactional; import jakarta.ws.rs.*; import jakarta.ws.rs.core.Context; import jakarta.ws.rs.core.MediaType; diff --git a/src/main/java/org/javastro/ivoa/tap/TAPJobService.java b/src/main/java/org/javastro/ivoa/tap/TAPJobService.java index d72c6dd..c8d9212 100644 --- a/src/main/java/org/javastro/ivoa/tap/TAPJobService.java +++ b/src/main/java/org/javastro/ivoa/tap/TAPJobService.java @@ -19,4 +19,9 @@ public TAPJobService(TAPHelper tapHelper) { public TAPJob createJob(TAPJobSpecification spec) throws UWSException { return (TAPJob) tapHelper.jobmanager.createJob(spec); } + + @Transactional + public boolean deleteJob(String jobId) throws UWSException { + return tapHelper.jobmanager.deleteJob(jobId); + } } From f2f74732a1d4324d77c95510452d6eefe6fa2bbf Mon Sep 17 00:00:00 2001 From: Stephen Lloyd Date: Mon, 29 Jun 2026 09:44:15 +0100 Subject: [PATCH 3/6] The delete job can be handled with a Transactional tag as it's not nested in a Uni --- .../java/org/javastro/ivoa/tap/TAPJobService.java | 14 +++++++++----- 1 file changed, 9 insertions(+), 5 deletions(-) diff --git a/src/main/java/org/javastro/ivoa/tap/TAPJobService.java b/src/main/java/org/javastro/ivoa/tap/TAPJobService.java index c8d9212..6e11b1d 100644 --- a/src/main/java/org/javastro/ivoa/tap/TAPJobService.java +++ b/src/main/java/org/javastro/ivoa/tap/TAPJobService.java @@ -6,6 +6,15 @@ import org.javastro.ivoacore.tap.TAPJobSpecification; import org.javastro.ivoacore.uws.UWSException; +/** + * TAPJobService is responsible for managing the creation of TAPJob instances + * based on provided specifications. It relies on TAPHelper to handle the + * underlying operations associated with job management. + *

+ * This service is application-scoped, ensuring a single instance is used + * throughout the application context. In a multithreaded environment, + * the service is designed to be thread-safe. + */ @ApplicationScoped public class TAPJobService { @@ -19,9 +28,4 @@ public TAPJobService(TAPHelper tapHelper) { public TAPJob createJob(TAPJobSpecification spec) throws UWSException { return (TAPJob) tapHelper.jobmanager.createJob(spec); } - - @Transactional - public boolean deleteJob(String jobId) throws UWSException { - return tapHelper.jobmanager.deleteJob(jobId); - } } From 44350fc6ac244f0ac8b1fffd2b4af15790e73c4e Mon Sep 17 00:00:00 2001 From: Stephen Lloyd Date: Mon, 29 Jun 2026 14:13:20 +0100 Subject: [PATCH 4/6] Move the sync service added for the cached job store (dbase transactions) due to the changes to hide most of the operations from the TapServer. Added the transactional overridden methods to the sync base class for the same reason. Sync approach requires a custom bean that can be indexed due to Uni moving the actual createJob call to a different thread --- .../quarkus/tap/BaseAsyncTAPResource.java | 42 ++++-- .../ivoa/quarkus/tap/BaseSyncTAPResource.java | 8 +- .../ivoa/quarkus}/tap/TAPJobService.java | 4 +- .../javastro/ivoa/tap/AsyncQueryResource.java | 74 +-------- .../org/javastro/ivoa/tap/QueryResource.java | 141 +----------------- src/main/resources/application.properties | 5 + 6 files changed, 52 insertions(+), 222 deletions(-) rename {src/main/java/org/javastro/ivoa => quarkus-tap-lib/src/main/java/org/javastro/ivoa/quarkus}/tap/TAPJobService.java (87%) diff --git a/quarkus-tap-lib/src/main/java/org/javastro/ivoa/quarkus/tap/BaseAsyncTAPResource.java b/quarkus-tap-lib/src/main/java/org/javastro/ivoa/quarkus/tap/BaseAsyncTAPResource.java index 3ee4155..9aea56f 100644 --- a/quarkus-tap-lib/src/main/java/org/javastro/ivoa/quarkus/tap/BaseAsyncTAPResource.java +++ b/quarkus-tap-lib/src/main/java/org/javastro/ivoa/quarkus/tap/BaseAsyncTAPResource.java @@ -6,6 +6,7 @@ package org.javastro.ivoa.quarkus.tap; +import jakarta.transaction.Transactional; import jakarta.ws.rs.*; import jakarta.ws.rs.core.*; import org.javastro.ivoa.quarkus.tap.upload.QuarkusTapUploader; @@ -46,18 +47,19 @@ protected Response redirectToJob(String jobid) { } //IMPL the two query endpoints are in different resources for routing purposes - @POST - @Consumes({MediaType.APPLICATION_FORM_URLENCODED, MediaType.MULTIPART_FORM_DATA}) - public Response async(@RestForm("QUERY") String query, @RestForm("LANG") String lang, @RestForm("RESPONSEFORMAT") String responseformat, - @RestForm("MAXREC") Long maxrec, @RestForm("RUNID") String runid, - @RestForm("UPLOAD") String upload, MultipartFormDataInput input, @Context UriInfo uriInfo) throws UWSException { - TAPUploadCacher tapUploader = new NullUploader(); - if(upload != null && !upload.isEmpty() ) { + @POST + @Consumes({MediaType.APPLICATION_FORM_URLENCODED, MediaType.MULTIPART_FORM_DATA}) + @Transactional + public Response async(@RestForm("QUERY") String query, @RestForm("LANG") String lang, @RestForm("RESPONSEFORMAT") String responseformat, + @RestForm("MAXREC") Long maxrec, @RestForm("RUNID") String runid, + @RestForm("UPLOAD") String upload, MultipartFormDataInput input, @Context UriInfo uriInfo) throws UWSException { + TAPUploadCacher tapUploader = new NullUploader(); + if(upload != null && !upload.isEmpty() ) { tapUploader = new QuarkusTapUploader(upload, input); - } - BaseUWSJob job = getTapHelper().getJobmanager().createJob(new TAPJobSpecification(query,lang,responseformat,maxrec,runid,tapUploader)); - return Response.seeOther(getTapHelper().asyncJobUri(job.getID())).build(); - } + } + BaseUWSJob job = getTapHelper().getJobmanager().createJob(new TAPJobSpecification(query,lang,responseformat,maxrec,runid,tapUploader)); + return Response.seeOther(getTapHelper().asyncJobUri(job.getID())).build(); + } @GET @Path("{jobid}/results/result") @@ -68,4 +70,22 @@ public RestResponse getVotable(@PathParam("jobid") String jo .header(HttpHeaders.CONTENT_DISPOSITION, "result.vot") .build(); } + +//----------------------- Need to make database modifying operations transactional ---------------------------------- +// Which means the base UWS modifying tasks need to be wrapped in a transactional override + @Override + @DELETE + @Path("{jobid}") + @Transactional + public Response deleteJob(@PathParam("jobid")String jobid) throws UWSException { + return super.deleteJob(jobid); + } + + @Override + @POST + @Path("{jobid}/phase") + @Transactional + public Response setPhase(@PathParam("jobid") String jobid, @FormParam("PHASE") String phase) throws UWSException { + return super.setPhase(jobid, phase); + } } diff --git a/quarkus-tap-lib/src/main/java/org/javastro/ivoa/quarkus/tap/BaseSyncTAPResource.java b/quarkus-tap-lib/src/main/java/org/javastro/ivoa/quarkus/tap/BaseSyncTAPResource.java index 5cae1e1..e8c02d6 100644 --- a/quarkus-tap-lib/src/main/java/org/javastro/ivoa/quarkus/tap/BaseSyncTAPResource.java +++ b/quarkus-tap-lib/src/main/java/org/javastro/ivoa/quarkus/tap/BaseSyncTAPResource.java @@ -8,6 +8,7 @@ import io.smallrye.mutiny.Uni; import io.smallrye.mutiny.infrastructure.Infrastructure; +import jakarta.inject.Inject; import jakarta.ws.rs.Consumes; import jakarta.ws.rs.GET; import jakarta.ws.rs.POST; @@ -43,6 +44,9 @@ */ public abstract class BaseSyncTAPResource { + @Inject + TAPJobService jobService; + protected static final Logger log = LoggerFactory.getLogger(BaseSyncTAPResource.class); /** @@ -91,9 +95,7 @@ protected Uni handleJob(String query, String lang, String responseformat, if(upload != null && !upload.isEmpty() ) { tapUploader = new QuarkusTapUploader(upload, input); } - job = (TAPJob) getTapHelper().getJobmanager().createJob( - new TAPJobSpecification(query, lang, responseformat, maxrec, runid, tapUploader) - ); + job = jobService.createJob(new TAPJobSpecification(query, lang, responseformat, maxrec, runid, tapUploader)); getTapHelper().getJobmanager().runJob(job.getID()); // automatically run the job } catch (UWSException e) { diff --git a/src/main/java/org/javastro/ivoa/tap/TAPJobService.java b/quarkus-tap-lib/src/main/java/org/javastro/ivoa/quarkus/tap/TAPJobService.java similarity index 87% rename from src/main/java/org/javastro/ivoa/tap/TAPJobService.java rename to quarkus-tap-lib/src/main/java/org/javastro/ivoa/quarkus/tap/TAPJobService.java index 6e11b1d..0c94953 100644 --- a/src/main/java/org/javastro/ivoa/tap/TAPJobService.java +++ b/quarkus-tap-lib/src/main/java/org/javastro/ivoa/quarkus/tap/TAPJobService.java @@ -1,4 +1,4 @@ -package org.javastro.ivoa.tap; +package org.javastro.ivoa.quarkus.tap; import jakarta.enterprise.context.ApplicationScoped; import jakarta.transaction.Transactional; @@ -26,6 +26,6 @@ public TAPJobService(TAPHelper tapHelper) { @Transactional public TAPJob createJob(TAPJobSpecification spec) throws UWSException { - return (TAPJob) tapHelper.jobmanager.createJob(spec); + return (TAPJob) tapHelper.getJobmanager().createJob(spec); } } diff --git a/src/main/java/org/javastro/ivoa/tap/AsyncQueryResource.java b/src/main/java/org/javastro/ivoa/tap/AsyncQueryResource.java index 2de6865..949f6e9 100644 --- a/src/main/java/org/javastro/ivoa/tap/AsyncQueryResource.java +++ b/src/main/java/org/javastro/ivoa/tap/AsyncQueryResource.java @@ -9,25 +9,10 @@ import jakarta.enterprise.context.ApplicationScoped; import jakarta.inject.Inject; import jakarta.ws.rs.*; -import jakarta.ws.rs.core.*; import org.eclipse.microprofile.openapi.annotations.tags.Tag; -import org.javastro.ivoa.tap.upload.QuarkusTapUploader; -import org.javastro.ivoacore.tap.TAPJobSpecification; -import org.javastro.ivoacore.tap.upload.NullUploader; -import org.javastro.ivoacore.tap.upload.TAPUploadCacher; -import org.javastro.ivoacore.uws.BaseUWSJob; -import org.javastro.ivoacore.uws.JobManager; -import org.javastro.ivoacore.uws.UWSException; -import org.javastro.ivoacore.uws.webapi.BaseUWSResource; -import org.jboss.resteasy.reactive.RestForm; -import org.jboss.resteasy.reactive.RestResponse; -import org.jboss.resteasy.reactive.server.multipart.MultipartFormDataInput; import org.javastro.ivoa.quarkus.tap.BaseAsyncTAPResource; import org.javastro.ivoa.quarkus.tap.TAPHelper; -import java.net.URI; -import java.util.Map; - /** * Main Async TAP Query. * Created on 04/03/2026 by Paul Harrison (paul.harrison@manchester.ac.uk). @@ -41,62 +26,7 @@ public class AsyncQueryResource extends BaseAsyncTAPResource { TAPHelper tapHelper; @Override - protected JobManager getJobManager() { - return tapHelper.jobmanager; - } - - @Override - protected Response redirectToJob(String jobid) { - - final UriBuilder urib = UriBuilder.fromUri(tapHelper.serviceLocator.serviceURI()) - .path("async"); - if (jobid != null && !jobid.isEmpty()) { - urib.path(jobid); - } - return Response.seeOther(urib - .build()).build(); - } - - //IMPL the two query endpoints are in different resources for routing purposes - @POST - @Consumes({MediaType.APPLICATION_FORM_URLENCODED, MediaType.MULTIPART_FORM_DATA}) - @Transactional - public Response async(@RestForm("QUERY") String query, @RestForm("LANG") String lang, @RestForm("RESPONSEFORMAT") String responseformat, - @RestForm("MAXREC") Long maxrec, @RestForm("RUNID") String runid, - @RestForm("UPLOAD") String upload, MultipartFormDataInput input, @Context UriInfo uriInfo) throws UWSException { - TAPUploadCacher tapUploader = new NullUploader(); - if(upload != null && !upload.isEmpty() ) { - tapUploader = new QuarkusTapUploader(upload, input); - } - BaseUWSJob job = tapHelper.jobmanager.createJob(new TAPJobSpecification(query,lang,responseformat,maxrec,runid,tapUploader)); - return Response.seeOther(tapHelper.asyncJobUri(job.getID())).build(); - } - - @GET - @Path("{jobid}/results/result") - @Produces("application/x-votable+xml") - public RestResponse getVotable(@PathParam("jobid") String jobid) throws UWSException { - final java.nio.file.Path path = tapHelper.getResultPath(jobid); - return RestResponse.ResponseBuilder.ok(path) - .header(HttpHeaders.CONTENT_DISPOSITION, "result.vot") - .build(); - } - - //----------------------- Need to make database modifying operations transactional ---------------------------------- - // Which means the base UWS modifying tasks need to be wrapped in a transactional override - @Override - @DELETE - @Path("{jobid}") - @Transactional - public Response deleteJob(@PathParam("jobid")String jobid) throws UWSException { - return super.deleteJob(jobid); - } - - @Override - @POST - @Path("{jobid}/phase") - @Transactional - public Response setPhase(@PathParam("jobid") String jobid, @FormParam("PHASE") String phase) throws UWSException { - return super.setPhase(jobid, phase); + protected TAPHelper getTapHelper() { + return tapHelper; } } diff --git a/src/main/java/org/javastro/ivoa/tap/QueryResource.java b/src/main/java/org/javastro/ivoa/tap/QueryResource.java index 036e44b..f783ecb 100644 --- a/src/main/java/org/javastro/ivoa/tap/QueryResource.java +++ b/src/main/java/org/javastro/ivoa/tap/QueryResource.java @@ -6,40 +6,14 @@ package org.javastro.ivoa.tap; -import io.smallrye.mutiny.Uni; -import io.smallrye.mutiny.infrastructure.Infrastructure; import jakarta.enterprise.context.ApplicationScoped; import jakarta.inject.Inject; import jakarta.ws.rs.*; -import jakarta.ws.rs.core.Context; -import jakarta.ws.rs.core.MediaType; -import jakarta.ws.rs.core.UriInfo; import org.eclipse.microprofile.config.inject.ConfigProperty; import org.eclipse.microprofile.openapi.annotations.tags.Tag; -import org.javastro.ivoa.entities.uws.ExecutionPhase; -import org.javastro.ivoa.tap.upload.QuarkusTapUploader; -import org.javastro.ivoacore.tap.TAPJob; -import org.javastro.ivoacore.tap.TAPJobSpecification; -import org.javastro.ivoacore.tap.TAPWriter; -import org.javastro.ivoacore.tap.upload.NullUploader; -import org.javastro.ivoacore.tap.upload.TAPUploadCacher; -import org.javastro.ivoacore.uws.UWSException; -import org.jboss.resteasy.reactive.RestForm; -import org.jboss.resteasy.reactive.RestQuery; -import org.jboss.resteasy.reactive.server.multipart.MultipartFormDataInput; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import uk.ac.starlink.table.*; import org.javastro.ivoa.quarkus.tap.BaseSyncTAPResource; import org.javastro.ivoa.quarkus.tap.TAPHelper; - -import java.io.IOException; -import java.net.URI; -import java.nio.file.Files; -import java.time.Duration; -import java.util.Map; - /** * Main TAP Query. * This works by creating an asynchronous job and then waiting for it to complete, @@ -52,120 +26,19 @@ @Path("sync") public class QueryResource extends BaseSyncTAPResource { - @ConfigProperty(name="ivoa.tap.sync-timeout-seconds", defaultValue = "5") int syncTimeoutSeconds; - @Inject - TAPJobService jobService; - - private static final Logger log = LoggerFactory.getLogger(QueryResource.class); - - @GET - @Produces("application/x-votable+xml") - public Uni syncGet(@RestQuery String query, @RestQuery String lang, @RestQuery String responseformat, @RestQuery Long maxrec, @RestQuery String runid, - @RestQuery String upload, - @Context UriInfo uriInfo) { - - return handleJob(query, lang, responseformat, maxrec, runid, upload, null, uriInfo); - } - - //UPLOAD param details - https://www.ivoa.net/documents/DALI/20170517/REC-DALI-1.1.html#tth_sEc3.4.5 - //UPLOAD=table1,http://example.com/t1.xml - //UPLOAD=image1,vos://example.authority!tempSpace/foo.fits - //UPLOAD=table3,param:t3 - @POST - @Consumes({MediaType.APPLICATION_FORM_URLENCODED, MediaType.MULTIPART_FORM_DATA}) - @Produces("application/x-votable+xml") - public Uni syncPost(@RestForm("QUERY") String query, @RestForm("LANG") String lang, @RestForm("RESPONSEFORMAT") String responseformat, @RestForm("MAXREC") Long maxrec, @RestForm("RUNID") String runid, - @RestForm("UPLOAD") String upload, - MultipartFormDataInput input, - @Context UriInfo uriInfo) { - - return handleJob(query, lang, responseformat, maxrec, runid, upload, input, uriInfo); - } - - private Uni handleJob(String query, String lang, String responseformat, Long maxrec, String runid, String upload, MultipartFormDataInput input, UriInfo uriInfo) { - final Duration SYNC_WAIT = Duration.ofSeconds(syncTimeoutSeconds); - return Uni.createFrom().deferred(() -> { - final TAPJob job; - try { - TAPUploadCacher tapUploader = new NullUploader(); - if(upload != null && !upload.isEmpty() ) { - tapUploader = new QuarkusTapUploader(upload, input); - } - // job = (TAPJob) tapHelper.jobmanager.createJob( - // new TAPJobSpecification(query, lang, responseformat, maxrec, runid, tapUploader) - // ); - job = jobService.createJob(new TAPJobSpecification(query, lang, responseformat, maxrec, runid, tapUploader)); - - tapHelper.jobmanager.runJob(job.getID()); // automatically run the job - } catch (UWSException e) { - return Uni.createFrom().failure(e); - } - - return Uni.createFrom().completionStage(job.getJobFuture()) - .onItem().transformToUni(phase -> { - if (phase == ExecutionPhase.COMPLETED) { - return successResponse(job); - } - else if (phase == ExecutionPhase.ERROR) - { - return Uni.createFrom().item( buildErrorVOTable(job,null, false)); - } - else { - return Uni.createFrom().failure(new UWSException("Underlying TAP job completed with unexpected phase " + phase));//TODO could do more sophisticated error handling here based on the phase - } - } - ) - .ifNoItem().after(SYNC_WAIT) - .recoverWithItem( - buildErrorVOTable(job, new UWSException("query did not complete within sync time limit of " + SYNC_WAIT.toSeconds() + " seconds - continuing as UWS job"), true)//FIXME should this return http error code - if so which code? - ); - - }).runSubscriptionOn(Infrastructure.getDefaultExecutor()); //TODO review whether this is the right way to do this - We might want to use a dedicated thread pool for this or some other strategy for managing the threads. - } + TAPHelper tapHelper; - private Uni successResponse(TAPJob job) { - return Uni.createFrom().item(() -> { - try { - return tapHelper.getResultPath(job.getID()); - } catch (UWSException e) { - throw new RuntimeException("Failed to get result path for job " + job.getID(), e); - } - }); + @Override + protected TAPHelper getTapHelper() { + return tapHelper; } - //TODO do we always want to return a VOTable even for errors? Or should we allow some other error response? - //TODO perhaps some of this can be moved to the TAPJob itself (for dealing with other types of errors - e.g. failure to parse original query) - protected java.nio.file.Path buildErrorVOTable(TAPJob job, UWSException exception, boolean timeout) { - // create a VOTable with STIL that has the error message and return the path to it. We could also include some info from the job if we have it. - - TAPJobSpecification tapJobSpec = (TAPJobSpecification) job.getJobSpecification(); - try { - - final TAPWriter tableWriter = new TAPWriter(job); - ColumnInfo[] columns = new ColumnInfo[]{ - new ColumnInfo("ERROR", String.class, "TAP error message") - }; - RowListStarTable table = new RowListStarTable(columns); - table.setName("error"); - if(exception != null) { - table.addRow(new Object[]{exception.getMessage()}); - } - if(timeout) { - tableWriter.setTimeoutInfo(tapHelper.asyncJobUri(job.getID())); - } - - java.nio.file.Path tempFile = java.nio.file.Files.createTempFile("error", ".vot"); - try (java.io.OutputStream out = java.nio.file.Files.newOutputStream(tempFile)) { - - tableWriter.writeStarTable(table, out); - } - return tempFile; - } catch (java.io.IOException e) { - throw new RuntimeException("Failed to create error VOTable: " + e.getMessage(), e); - } + @Override + protected int getSyncWait() { + return syncTimeoutSeconds; } } diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties index 7e02119..edbd94a 100644 --- a/src/main/resources/application.properties +++ b/src/main/resources/application.properties @@ -70,3 +70,8 @@ quarkus.kubernetes.ingress.path-type=Prefix #attempt to reduce logging noise from netty websocket quarkus.log.category."io.netty.handler.codec.http.websocketx".use-parent-handlers=false + +# quarkus-tap-lib contains base JAX-RS/CDI resource classes with interceptor annotations +# such as @Transactional. Index it so Quarkus can see those annotations at build time. +quarkus.index-dependency.quarkus-tap-lib.group-id=org.javastro.ivoa.servers +quarkus.index-dependency.quarkus-tap-lib.artifact-id=quarkus-tap-lib From 62af3d3c8818680ac5f369348ad74b1218740397 Mon Sep 17 00:00:00 2001 From: Paul Harrison Date: Wed, 1 Jul 2026 14:22:09 +0100 Subject: [PATCH 5/6] factor out the uwsstore config to be distinct from the tap store --- .../javastro/ivoa/tap/TapConfiguration.java | 2 ++ src/main/resources/application.properties | 24 +++++++++++-------- 2 files changed, 16 insertions(+), 10 deletions(-) diff --git a/src/main/java/org/javastro/ivoa/tap/TapConfiguration.java b/src/main/java/org/javastro/ivoa/tap/TapConfiguration.java index 2d18b79..748677e 100644 --- a/src/main/java/org/javastro/ivoa/tap/TapConfiguration.java +++ b/src/main/java/org/javastro/ivoa/tap/TapConfiguration.java @@ -6,6 +6,7 @@ package org.javastro.ivoa.tap; +import io.quarkus.hibernate.orm.PersistenceUnit; import jakarta.enterprise.context.ApplicationScoped; import jakarta.enterprise.inject.Produces; import jakarta.inject.Inject; @@ -53,6 +54,7 @@ public class TapConfiguration { DataSource ds; @Inject + @PersistenceUnit("uwsstore") EntityManager em; @ConfigProperty(name="ivoa.tap.dbCaseSensitive", defaultValue = "false") diff --git a/src/main/resources/application.properties b/src/main/resources/application.properties index edbd94a..bab6dcd 100644 --- a/src/main/resources/application.properties +++ b/src/main/resources/application.properties @@ -23,30 +23,34 @@ quarkus.datasource.devservices.enabled=true #fix port %dev.quarkus.datasource.devservices.port=54771 quarkus.datasource.devservices.image-name=images.dev.uksrc.org/library/postgres-pgsphere14:latest -#quarkus.datasource.devservices.image-name=postgis/postgis:14-3.5 # # PGSphere quarkus.hibernate-orm.dialect=org.javastro.ivoacore.pgsphere.PgSphereDialect +quarkus.hibernate-orm.schema-management.create-schemas=true -#ORM setup - mainly for the TAPSchema -quarkus.hibernate-orm.packages=org.ivoa.dm,org.ivoa.vodml.stdtypes,org.javastro.ivoacore.pgsphere,org.javastro.ivoa.tap.entities,org.javastro.ivoacore.uws.persist -#%dev,%test.quarkus.hibernate-orm.database.generation=drop-and-create -%prod.quarkus.hibernate-orm.database.generation=update -quarkus.hibernate-orm.database.generation.create-schemas=true -quarkus.hibernate-orm.database.generation.halt-on-error=false +#default PU ORM setup - mainly for the TAPSchema +quarkus.hibernate-orm.packages=org.ivoa.dm,org.ivoa.vodml.stdtypes,org.javastro.ivoacore.pgsphere,org.javastro.ivoa.tap.entities quarkus.hibernate-orm.quote-identifiers.strategy = all # below is new quarkus -#quarkus.hibernate-orm.schema-management.strategy=drop-and-create -# older quarkus +quarkus.hibernate-orm.schema-management.strategy=drop-and-create +# this is for the scripts quarkus.hibernate-orm.scripts.generation=drop-and-create quarkus.hibernate-orm.scripts.generation.create-target=TAPddl.sql quarkus.hibernate-orm.scripts.generation.drop-target=TAP_drop_ddl.sql - quarkus.hibernate-orm.log.sql=true quarkus.hibernate-orm.log.bind-parameters=true +#datasource/PU ORM setup for UWS storage (TODO are the quotes necessary? The examples have them) +#if a separate database is needed then the datasource should be created and the datasource name should be used here instead of +#quarkus.datasource."uwsstore".db-kind=postgresql +quarkus.hibernate-orm."uwsstore".packages=org.javastro.ivoacore.uws.persist +quarkus.hibernate-orm."uwsstore".schema-management.strategy=drop-and-create +quarkus.hibernate-orm."uwsstore".schema-management.create-schemas=true +quarkus.hibernate-orm."uwsstore".datasource= + + #image build quarkus.container-image.builder=docker From 2bb95b95c5f8f9d22edb2670e03a74b93d8a70bd Mon Sep 17 00:00:00 2001 From: Paul Harrison Date: Thu, 2 Jul 2026 15:58:13 +0100 Subject: [PATCH 6/6] simplified interface to the cached store. --- .../java/org/javastro/ivoa/tap/TapConfiguration.java | 10 ++-------- 1 file changed, 2 insertions(+), 8 deletions(-) diff --git a/src/main/java/org/javastro/ivoa/tap/TapConfiguration.java b/src/main/java/org/javastro/ivoa/tap/TapConfiguration.java index 748677e..c667b22 100644 --- a/src/main/java/org/javastro/ivoa/tap/TapConfiguration.java +++ b/src/main/java/org/javastro/ivoa/tap/TapConfiguration.java @@ -114,24 +114,18 @@ JobManager uws(SchemaProvider schemaProvider) { } EnvironmentFactory env = new DefaultEnvironmentFactory(tmpdir); - // MemoryBasedJobStore store = new MemoryBasedJobStore(); TAPJob.JobFactory tapJobFactory = new TAPJob.JobFactory(ds, schemaProvider, env); - JobFactoryAggregator factoryAggregator = new JobFactoryAggregator(); - factoryAggregator.addFactory(tapJobFactory); - JobStore store = new CachedJobStore( DatabaseJobStore.forJobType( em, TAPJobSpecification.class, - "TAP", - factoryAggregator + "TAP" ) ); - //CachedJobStore store = new CachedJobStore(new DatabaseJobStore(ds, em, )); DefaultExecutionPolicy policy = new DefaultExecutionPolicy(); - return new JobManager(factoryAggregator, store, policy); + return new JobManager(tapJobFactory, store, policy); } @Produces