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
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import au.org.aodn.ogcapi.server.core.util.RestTemplateUtils;
import au.org.aodn.ogcapi.server.processes.BatchJobProperties;
import au.org.aodn.ogcapi.server.processes.DownloadLimitProperties;
import au.org.aodn.ogcapi.server.processes.DownloadShareProperties;
import au.org.aodn.ogcapi.server.processes.DownloadSizeLimitProperties;
import com.fasterxml.jackson.annotation.JsonInclude;
import com.fasterxml.jackson.databind.ObjectMapper;
Expand Down Expand Up @@ -39,6 +40,7 @@
BatchJobProperties.class,
DownloadLimitProperties.class,
DownloadSizeLimitProperties.class,
DownloadShareProperties.class,
OgcApiProperties.class
})
public class Config {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,8 @@ public enum Parameter {
FULL_METADATA_LINK("full_metadata_link"),
SUGGESTED_CITATION("suggested_citation"),
KEY("key"),
OUTPUT_FORMAT("output_format");
OUTPUT_FORMAT("output_format"),
SHARE_IDENTIFIER("share_identifier");

private final String value;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,13 +2,15 @@

import au.org.aodn.ogcapi.server.core.exception.DownloadLimitExceededException;
import au.org.aodn.ogcapi.server.core.exception.DownloadSizeExceededException;
import au.org.aodn.ogcapi.server.core.model.enumeration.DatasetDownloadEnums;
import com.fasterxml.jackson.core.JsonProcessingException;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;

import java.util.HashMap;
import java.util.Map;
import java.util.OptionalLong;
import java.util.concurrent.locks.ReentrantLock;

/**
Expand All @@ -17,6 +19,7 @@
* the subset.
* 2. Concurrency: a user already at the limit is rejected outright, the caller has to wait for
* one of their own downloads to finish and try again.
* An admitted download is tagged with its fair-share (small or large) from the same estimate.
*/
@Slf4j
@Service
Expand All @@ -26,6 +29,7 @@ public class DownloadAdmissionService {
private final InFlightDownloadCounter counter;
private final DownloadLimitProperties limits;
private final DownloadSizeLimitProperties sizeLimit;
private final DownloadShareProperties share;

private final ReentrantLock lock = new ReentrantLock();

Expand All @@ -43,11 +47,13 @@ public DownloadAdmissionService(
RestServices restServices,
InFlightDownloadCounter counter,
DownloadLimitProperties limits,
DownloadSizeLimitProperties sizeLimit) {
DownloadSizeLimitProperties sizeLimit,
DownloadShareProperties share) {
this.restServices = restServices;
this.counter = counter;
this.limits = limits;
this.sizeLimit = sizeLimit;
this.share = share;
}

/**
Expand All @@ -60,10 +66,8 @@ public DownloadAdmissionService(
public String submit(DownloadRequest request) throws JsonProcessingException {
String key = InFlightDownloadCounter.recipientKey(request.recipient());

if (sizeLimit.enabled()) {
// Before the slot is reserved, so a slow estimate does not hold one.
rejectIfTooLarge(request);
}
// Before the slot is reserved, so a slow estimate does not hold one.
OptionalLong estimatedBytes = sizeLimit.enabled() ? rejectIfTooLarge(request) : OptionalLong.empty();

if (limits.enabled()) {
// Outside the lock: the sweep is the only expensive step.
Expand All @@ -73,8 +77,14 @@ public String submit(DownloadRequest request) throws JsonProcessingException {

try {
Map<String, String> parameters = restServices.buildDownloadParameters(request);
String shareIdentifier = shareFor(estimatedBytes);
log.info("Download for uuid {} estimated {} bytes, share {}", request.uuid(),
estimatedBytes.isPresent() ? estimatedBytes.getAsLong() : "none", shareIdentifier);
// The Batch shareIdentifier is only sent where the queue has a fair-share policy, otherwise AWS rejects the submit.
parameters.put(DatasetDownloadEnums.Parameter.SHARE_IDENTIFIER.getValue(), shareIdentifier);
String jobName = RestServices.downloadJobName(request.recipient());
String awsJobId = restServices.submitDownloadJob(jobName, parameters);
String awsJobId = restServices.submitDownloadJob(
jobName, parameters, share.enabled() ? shareIdentifier : null);
counter.recordSubmitted(awsJobId, request.recipient());
notifyStarted(request);
return awsJobId;
Expand All @@ -85,15 +95,18 @@ public String submit(DownloadRequest request) throws JsonProcessingException {
}
}

private void rejectIfTooLarge(DownloadRequest request) {
/**
* @return the estimate, or empty when DAS could not give one
*/
private OptionalLong rejectIfTooLarge(DownloadRequest request) {
long estimatedBytes;
try {
estimatedBytes = restServices.estimateDownloadBytes(request);
} catch (Exception e) {
// Let it through, as the portal does when its own estimate fails. Blocking every
// download because DAS cannot estimate would be worse than one job running out of disk.
log.warn("Size estimate failed for uuid {}, submitting without the size check", request.uuid(), e);
return;
return OptionalLong.empty();
}

long maxBytes = sizeLimit.maxSize().toBytes();
Expand All @@ -102,6 +115,18 @@ private void rejectIfTooLarge(DownloadRequest request) {
request.uuid(), estimatedBytes, maxBytes);
throw new DownloadSizeExceededException(estimatedBytes, maxBytes);
}
return OptionalLong.of(estimatedBytes);
}

/**
* The fair-share tag for a download: small only when it is estimated under the threshold.
* No estimate counts as large, so a download we know nothing about cannot take the slots
* kept for small ones.
*/
private String shareFor(OptionalLong estimatedBytes) {
boolean small = estimatedBytes.isPresent()
&& estimatedBytes.getAsLong() < share.smallMaxSize().toBytes();
return small ? share.small() : share.large();
}

private void reserveOrReject(DownloadRequest request, String key) {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
package au.org.aodn.ogcapi.server.processes;

import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.boot.context.properties.bind.DefaultValue;
import org.springframework.util.unit.DataSize;

/**
* Which fair-share a download belongs to. A download estimated under smallMaxSize is tagged
* small; anything else, including a download with no estimate, is tagged large. The tag always
* goes out as the share_identifier job parameter. When enabled is true it is also sent as the
* Batch shareIdentifier, which allows only alphanumerics. Turn it on only where the queue has a
* fair-share scheduling policy, because AWS rejects the field on a FIFO queue and requires it on
* a fair-share one.
*/
@ConfigurationProperties(prefix = "aws.batch.job.share")
public record DownloadShareProperties(
@DefaultValue("false") boolean enabled,
@DefaultValue("50MB") DataSize smallMaxSize,
@DefaultValue("small") String small,
@DefaultValue("large") String large
) {
}
Original file line number Diff line number Diff line change
Expand Up @@ -151,14 +151,16 @@ public static String downloadJobName(String recipient) {

/**
* Submit a prepared download to the configured queue and job definition.
*
* @param shareIdentifier the fair-share tag sent as the Batch shareIdentifier, default as false, set in application.yaml
*/
public String submitDownloadJob(String jobName, Map<String, String> parameters) {
String jobId = submitJob(jobName, this.batchJobQueue, this.batchJobDefinition, parameters);
public String submitDownloadJob(String jobName, Map<String, String> parameters, String shareIdentifier) {
String jobId = submitJob(jobName, this.batchJobQueue, this.batchJobDefinition, parameters, shareIdentifier);
log.info("Job submitted with ID: {}", jobId);
return jobId;
}

private String submitJob(String jobName, String jobQueue, String jobDefinition, Map<String, String> parameters) {
private String submitJob(String jobName, String jobQueue, String jobDefinition, Map<String, String> parameters, String shareIdentifier) {

// Filter out null or empty parameter values before submitting to AWS Batch.
// AWS Batch returns "Parameter values must be provided" when the job definition
Expand All @@ -183,6 +185,8 @@ private String submitJob(String jobName, String jobQueue, String jobDefinition,
.jobQueue(jobQueue)
.jobDefinition(jobDefinition)
.parameters(submitParameters)
// SDK v2 leaves a null field out of the request, so FIFO submits are unchanged.
.shareIdentifier(shareIdentifier)
.build();

SubmitJobResponse submitJobResponse = batchClient.submitJob(submitJobRequest);
Expand Down
11 changes: 11 additions & 0 deletions server/src/main/resources/application.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,17 @@ aws:
size-limit:
enabled: true
max-size: 180GB
# Fair-share tag for a download (#9348). A download estimated under small-max-size is
# tagged small, anything else (including no estimate) large. The tag always goes out as
# the share_identifier job parameter. enabled also sets it as the Batch shareIdentifier;
# set it true only for an environment whose queue has a fair-share policy, and repoint
# both queue and child-queue there. child-queue is set explicitly above, so overriding
# queue alone leaves child-job lookups on the old queue.
share:
enabled: false
small-max-size: 50MB
small: small
large: large

wfs-default-param:
fields:
Expand Down
Loading
Loading