Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 16 additions & 9 deletions src/main/java/org/javastro/ivoa/tap/AsyncQueryResource.java
Original file line number Diff line number Diff line change
Expand Up @@ -11,14 +11,20 @@
import jakarta.ws.rs.*;
import jakarta.ws.rs.core.*;
import org.eclipse.microprofile.openapi.annotations.tags.Tag;
import org.javastro.ivoacore.common.ServiceLocator;
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 java.net.URI;
import java.util.Map;

/**
* Main Async TAP Query.
Expand All @@ -29,8 +35,7 @@
@Path("async")
public class AsyncQueryResource extends BaseUWSResource {


@Inject
@Inject
TAPHelper tapHelper;

@Override
Expand All @@ -50,13 +55,17 @@ protected Response redirectToJob(String jobid) {
.build()).build();
}


//IMPL the two query endpoints are in different resources for routing purposes.
//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, @Context UriInfo uriInfo) throws UWSException {
BaseUWSJob job = tapHelper.jobmanager.createJob(new TAPJobSpecification(query,lang,responseformat,maxrec,runid,upload));
@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();
}

Expand All @@ -69,6 +78,4 @@ public RestResponse<java.nio.file.Path> getVotable(@PathParam("jobid") String jo
.header(HttpHeaders.CONTENT_DISPOSITION, "result.vot")
.build();
}


}
43 changes: 31 additions & 12 deletions src/main/java/org/javastro/ivoa/tap/QueryResource.java
Original file line number Diff line number Diff line change
Expand Up @@ -10,25 +10,33 @@
import io.smallrye.mutiny.infrastructure.Infrastructure;
import jakarta.enterprise.context.ApplicationScoped;
import jakarta.inject.Inject;
import jakarta.ws.rs.GET;
import jakarta.ws.rs.POST;
import jakarta.ws.rs.Path;
import jakarta.ws.rs.Produces;
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 uk.ac.starlink.table.ColumnInfo;
import uk.ac.starlink.table.RowListStarTable;
import org.jboss.resteasy.reactive.server.multipart.MultipartFormDataInput;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
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.
Expand All @@ -48,31 +56,44 @@ public class QueryResource {
@Inject
TAPHelper tapHelper;

private static final Logger log = LoggerFactory.getLogger(QueryResource.class);

@GET
@Produces("application/x-votable+xml")
public Uni<java.nio.file.Path> 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, 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<java.nio.file.Path> 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, uriInfo);

return handleJob(query, lang, responseformat, maxrec, runid, upload, input, uriInfo);
}


private Uni<java.nio.file.Path> handleJob(String query, String lang, String responseformat, Long maxrec, String runid, String upload, UriInfo uriInfo) {
private Uni<java.nio.file.Path> 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, upload)
new TAPJobSpecification(query, lang, responseformat, maxrec, runid, tapUploader)
);

tapHelper.jobmanager.runJob(job.getID()); // automatically run the job
Expand Down Expand Up @@ -143,6 +164,4 @@ protected java.nio.file.Path buildErrorVOTable(TAPJob job, UWSException exceptio
throw new RuntimeException("Failed to create error VOTable: " + e.getMessage(), e);
}
}


}
64 changes: 64 additions & 0 deletions src/main/java/org/javastro/ivoa/tap/upload/QuarkusTapUploader.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
package org.javastro.ivoa.tap.upload;

import org.javastro.ivoacore.tap.upload.BaseTAPUploadCacher;
import org.javastro.ivoacore.tap.upload.TapUploadService;
import org.jboss.resteasy.reactive.server.multipart.FormValue;
import org.jboss.resteasy.reactive.server.multipart.MultipartFormDataInput;
import org.jspecify.annotations.NonNull;

import java.io.IOException;
import java.net.URI;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.StandardCopyOption;
import java.util.Map;
import java.util.Optional;
import java.util.UUID;

/**
* The QuarkusTapUploader class provides functionality for parsing and processing
* the DALI-compliant UPLOAD parameter, handling file uploads or URLs for
* input data. It includes methods for extracting upload specifications
* from the UPLOAD parameter, storing VOTables as temporary files, and
* generating appropriate Paths.
*/
public class QuarkusTapUploader extends BaseTAPUploadCacher {

private final MultipartFormDataInput input;

public QuarkusTapUploader(String uploadParam, MultipartFormDataInput input) {
super(uploadParam);
this.input = input;
}

/**
* Stores a VOTable in a temporary file and returns the URI of the file.
* @param theParam a param parameter of the DALI UPLOAD query parameter, e.g. "param:t3"
* @return The Path of the uploaded file, or null if the file was not uploaded.
* @throws IOException If an I/O error occurs while storing the file.
*/
protected Path storeParam(@NonNull String theParam, Path dir) throws IOException {
String paramName = theParam.split(":")[1];

Optional<FormValue> value = Optional.ofNullable(input.getValues().get(paramName))
.flatMap(list -> list.stream().findFirst());

if (value.isPresent() && value.get().isFileItem()) {
java.nio.file.Path uploadedFile = value.get().getFileItem().getFile();
Path persistent = generateFileName(dir, paramName);
Files.copy(uploadedFile, persistent, StandardCopyOption.REPLACE_EXISTING);

return persistent;
}
return null;
}

@Override
public boolean hasUpload() {
return true;
}




}
85 changes: 85 additions & 0 deletions src/test/java/org/javastro/ivoa/tap/AbstractTAPTest.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,85 @@
/*
* Copyright (c) 2026. Paul Harrison, University of Manchester.
*
*/

package org.javastro.ivoa.tap;


import io.restassured.filter.log.LogDetail;
import io.restassured.http.ContentType;

import java.time.Duration;

import static io.restassured.RestAssured.given;
import static org.awaitility.Awaitility.await;
import static org.hamcrest.Matchers.*;
import static org.junit.jupiter.api.Assertions.fail;

public abstract class AbstractTAPTest {
protected static void startAndTestJob(String jobUrl) {
// Start the Job (POST to /phase with PHASE=RUN)
given()
.contentType(ContentType.URLENC)
.redirects().follow(false)
.formParam("PHASE", "RUN")
.when()
.post(jobUrl + "/phase")
.then()
.statusCode(303);

// Poll until
await().atMost(Duration.ofSeconds(30))
.pollInterval(Duration.ofMillis(500))
.untilAsserted(() -> {
given()
.when().get(jobUrl + "/phase")
.then()
.statusCode(200)
.body(not(comparesEqualTo("RUNNING")));
})
;
String status = given()
.when().get(jobUrl + "/phase")
.then()
.statusCode(200)
.extract().body().asString();
if (status.equals("ERROR"))
{

given()
.when().get(jobUrl + "/error")
.then()
.statusCode(200)
.log().body();
fail("Job ended in error state");

}
else if (status.equals("COMPLETED")) {


// Retrieve Results
given()
.when().get(jobUrl + "/results")
.then()
.log().ifValidationFails(LogDetail.BODY)
.statusCode(200)
.body("results.result.size()", greaterThan(0));

//retrieve the actual result
given()
.when().get(jobUrl + "/results/result")
.then()
.statusCode(200)
.log().body();
;

//TODO get the result into a file and verify that it is an OK VOTable.

}
else
{
fail("Unexpected job status "+status);
}
}
}
Loading
Loading