diff --git a/adapters/klab.ogc/src/main/java/org/integratedmodelling/klab/stac/STACEncoder.java b/adapters/klab.ogc/src/main/java/org/integratedmodelling/klab/stac/STACEncoder.java index 53e597351..09006d1c9 100644 --- a/adapters/klab.ogc/src/main/java/org/integratedmodelling/klab/stac/STACEncoder.java +++ b/adapters/klab.ogc/src/main/java/org/integratedmodelling/klab/stac/STACEncoder.java @@ -2,9 +2,11 @@ import java.io.IOException; import java.io.OutputStream; +import java.net.URL; import java.time.LocalDateTime; import java.time.ZoneOffset; import java.time.format.DateTimeFormatter; +import java.util.ArrayList; import java.util.Date; import java.util.List; import java.util.Map; @@ -63,23 +65,30 @@ import org.integratedmodelling.klab.exceptions.KlabInternalErrorException; import org.integratedmodelling.klab.exceptions.KlabResourceAccessException; import org.integratedmodelling.klab.exceptions.KlabValidationException; +import org.integratedmodelling.klab.ogc.OpenEOAdapter; import org.integratedmodelling.klab.ogc.STACAdapter; import org.integratedmodelling.klab.ogc.vector.files.VectorEncoder; +import org.integratedmodelling.klab.openeo.OpenEO; +import org.integratedmodelling.klab.openeo.OpenEO.OpenEOFuture; +import org.integratedmodelling.klab.openeo.OpenEO.Process; import org.integratedmodelling.klab.raster.files.RasterEncoder; import org.integratedmodelling.klab.rest.ExternalAuthenticationCredentials; import org.integratedmodelling.klab.scale.Scale; import org.integratedmodelling.klab.stac.extensions.COGAssetExtension; import org.integratedmodelling.klab.stac.extensions.STACFeatureExtension; import org.integratedmodelling.klab.stac.extensions.STACIIASAExtension; +import org.integratedmodelling.klab.utils.JsonUtils; import org.integratedmodelling.klab.utils.s3.S3URLUtils; import org.locationtech.jts.geom.Envelope; import org.locationtech.jts.geom.Geometry; import org.locationtech.jts.geom.Polygon; +import org.locationtech.jts.operation.union.UnaryUnionOp; import com.fasterxml.jackson.databind.node.ArrayNode; import com.github.davidmoten.aws.lw.client.Client; import com.github.davidmoten.aws.lw.client.Credentials; +import kong.unirest.json.JSONArray; import kong.unirest.json.JSONObject; public class STACEncoder implements IResourceEncoder { @@ -360,7 +369,7 @@ public boolean test(HMStacAsset asset) { // Assuming for now that "eo:bands" wou HMStacCollection collection = null; try { manager.open(); - collection = manager.getCollectionById(resource.getParameters().get("collectionId", String.class)); + collection = manager.getCollectionById(collectionId); } catch (Exception e1) { throw new KlabResourceAccessException("Cannot access to STAC collection " + collectionUrl); } @@ -405,7 +414,7 @@ public boolean test(HMStacAsset asset) { // Assuming for now that "eo:bands" wou HMStacCollection collection = null; try { manager.open(); - collection = manager.getCollectionById(resource.getParameters().get("collectionId", String.class)); + collection = manager.getCollectionById(collectionId); if (collection == null) { scope.getMonitor() @@ -446,6 +455,7 @@ public boolean test(HMStacAsset asset) { // Assuming for now that "eo:bands" wou manager.close(); return; } + List items = collection.searchItems(); if (items.isEmpty()) { @@ -499,6 +509,57 @@ public boolean test(HMStacAsset asset) { // Assuming for now that "eo:bands" wou scope.getMonitor().warn("Multiple EPSGs found on the assets in items " + EPSGAtAssets.toString() + "." + "The transformation process could affect the data."); } + + + // Specific Implementation for the Slow Requests flow in WEED + if (collection.getId().contains("EU_modelV2-1-MECE") + && resource.getUrn().contains("im.resources-main")) { + Geometry unionMLStacInference = UnaryUnionOp.union( + items.stream() + .map(item -> item.getGeometry()) + .filter(g -> g != null && !g.isEmpty()) + .collect(Collectors.toList()) + ); + + if (!unionMLStacInference.contains(poly)) { + scope.getMonitor().warn("The requested extend for ML inferences is not completely contained in STAC, Starting ML Inference Request"); + + OpenEO service = OpenEOAdapter.getClient("openeo_weed.dataspace.copernicus.eu"); + List processes = new ArrayList<>(); + String processNamespace = "https://raw.githubusercontent.com/ESA-WEED-project/OpenEO-UDP-UDF-catalogue/refs/heads/main/UDP/json/udp_starter.json"; + String processID = "udp_starter"; + + Process process = JsonUtils.load(new URL(processNamespace), + Process.class); + process.encodeSelf(processNamespace); + processes.add(process); + + JSONObject arguments = new JSONObject() + .put("bbox", new JSONObject() + .put("crs", 4326) + .put("west", bbox.get(0)) + .put("south", bbox.get(3)) + .put("east", bbox.get(1)) + .put("north", bbox.get(2))) + .put("digitalId", "AM1729") // Forms the STAC coordinate later + .put("scenarioId", "DT_SLOW_FLOW") // Forms the STAC coordinate later + .put("year", ctxTime.getEnd().getYear()) + .put("onnx_model", "EUNIS2021plus_panEU_v201_2024_OneZone") // Hardcoding for now only for Europe, until the "BEST" model is decided! + .put("dt_url", "https://services.integratedmodelling.org/runtime/main/api/v1/dt/ESA_INSTITUTIONAL.3vh554o6h6c"); + + + OpenEOFuture job = service.submit(processID, arguments, + scope.getMonitor(), processes.toArray(new Process[processes.size()])); + + if (job.isCancelled()) { + scope.getMonitor().warn("job canceled"); + } else if (job.getError() != null) { + scope.getMonitor().error(job.getError()); + } else { + scope.getMonitor().info("Inference Request has been submitted to the ML Workflows"); + } + } + } HMRaster outRaster = collection.readRasterBandOnRegion(regionTransformed, assetPredicate, items, allowTransform, MergeMode.SUBSTITUTE, lpm); diff --git a/adapters/klab.ogc/src/main/java/org/integratedmodelling/klab/stac/extensions/STACFeatureExtension.java b/adapters/klab.ogc/src/main/java/org/integratedmodelling/klab/stac/extensions/STACFeatureExtension.java index aa668f566..b422e8122 100644 --- a/adapters/klab.ogc/src/main/java/org/integratedmodelling/klab/stac/extensions/STACFeatureExtension.java +++ b/adapters/klab.ogc/src/main/java/org/integratedmodelling/klab/stac/extensions/STACFeatureExtension.java @@ -42,7 +42,7 @@ public static FeatureSource getFeatures(JSONOb JSONObject searchPayload = new JSONObject() - .put("limit", 100) + .put("limit", 1000) .put("bbox", bboxArray) .put("collections", new JSONArray().put(collectionId)); @@ -66,26 +66,25 @@ public static FeatureSource getFeatures(JSONOb HMStacItem item = HMStacItem.fromSimpleFeature(feat); DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"); - if (item.getStartTimestamp() != null && item.getEndTimestamp() != null) { - long itemStart = LocalDateTime - .parse(item.getStartTimestamp(), formatter) - .atZone(ZoneOffset.UTC) - .toInstant() - .toEpochMilli(); - - long itemEnd = LocalDateTime - .parse(item.getEndTimestamp(), formatter) - .atZone(ZoneOffset.UTC) - .toInstant() - .toEpochMilli(); - if (start.getMilliseconds() >= itemStart && end.getMilliseconds() <= itemEnd) { - featureList.add(feat); - } - } else { - featureList.add(feat); - } - + if (item.getStartTimestamp() == null || item.getEndTimestamp() == null) { // Assume best case scenario + featureList.add(feat); + continue; + } + + long itemStart = LocalDateTime + .parse(item.getStartTimestamp(), formatter) + .atZone(ZoneOffset.UTC) + .toInstant() + .toEpochMilli(); + long itemEnd = LocalDateTime + .parse(item.getEndTimestamp(), formatter) + .atZone(ZoneOffset.UTC) + .toInstant() + .toEpochMilli(); + if (start.getMilliseconds() >= itemStart && end.getMilliseconds() <= itemEnd) { + featureList.add(feat); + } } catch (Exception e) { e.printStackTrace(); } diff --git a/klab.authentication/src/main/java/org/integratedmodelling/klab/Authentication.java b/klab.authentication/src/main/java/org/integratedmodelling/klab/Authentication.java index 2d980bbdf..9c7761261 100644 --- a/klab.authentication/src/main/java/org/integratedmodelling/klab/Authentication.java +++ b/klab.authentication/src/main/java/org/integratedmodelling/klab/Authentication.java @@ -463,7 +463,7 @@ public ExternalAuthenticationCredentials getCredentials(String endpoint) { * @return the credential or null if not present */ private ExternalAuthenticationCredentials getCredentialsByUrl(String hostUrl) { - return externalCredentials.values().stream().filter(credential -> credential.getURL().equals(hostUrl)).sorted() + return externalCredentials.values().stream().filter(credential -> credential.getURL().equals(hostUrl)) .findFirst().orElse(null); }