From 350fc911c171ec8b38ac2bfdb90779fff9150cdb Mon Sep 17 00:00:00 2001 From: Yuxuan HU Date: Wed, 7 Oct 2026 15:31:51 +1100 Subject: [PATCH 1/7] add shared share_identifier and properties --- .../enumeration/DatasetDownloadEnums.java | 3 ++- .../processes/DownloadShareProperties.java | 18 ++++++++++++++++++ 2 files changed, 20 insertions(+), 1 deletion(-) create mode 100644 server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadShareProperties.java diff --git a/server/src/main/java/au/org/aodn/ogcapi/server/core/model/enumeration/DatasetDownloadEnums.java b/server/src/main/java/au/org/aodn/ogcapi/server/core/model/enumeration/DatasetDownloadEnums.java index e272fd62..cd423ccd 100644 --- a/server/src/main/java/au/org/aodn/ogcapi/server/core/model/enumeration/DatasetDownloadEnums.java +++ b/server/src/main/java/au/org/aodn/ogcapi/server/core/model/enumeration/DatasetDownloadEnums.java @@ -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; diff --git a/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadShareProperties.java b/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadShareProperties.java new file mode 100644 index 00000000..c744bb42 --- /dev/null +++ b/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadShareProperties.java @@ -0,0 +1,18 @@ +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 values go + * into the Batch shareIdentifier later, which allows only alphanumerics. + */ +@ConfigurationProperties(prefix = "aws.batch.job.share") +public record DownloadShareProperties( + @DefaultValue("50MB") DataSize smallMaxSize, + @DefaultValue("small") String small, + @DefaultValue("large") String large +) { +} From 3f30b5f0b19dc327a12944c96fc84a03d4c48541 Mon Sep 17 00:00:00 2001 From: Yuxuan HU Date: Wed, 7 Oct 2026 15:33:58 +1100 Subject: [PATCH 2/7] add shared share_identifier config --- .../org/aodn/ogcapi/server/core/configuration/Config.java | 2 ++ server/src/main/resources/application.yaml | 7 +++++++ 2 files changed, 9 insertions(+) diff --git a/server/src/main/java/au/org/aodn/ogcapi/server/core/configuration/Config.java b/server/src/main/java/au/org/aodn/ogcapi/server/core/configuration/Config.java index 6cbbc629..a575b665 100644 --- a/server/src/main/java/au/org/aodn/ogcapi/server/core/configuration/Config.java +++ b/server/src/main/java/au/org/aodn/ogcapi/server/core/configuration/Config.java @@ -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; @@ -39,6 +40,7 @@ BatchJobProperties.class, DownloadLimitProperties.class, DownloadSizeLimitProperties.class, + DownloadShareProperties.class, OgcApiProperties.class }) public class Config { diff --git a/server/src/main/resources/application.yaml b/server/src/main/resources/application.yaml index 9238a48a..dab0baad 100644 --- a/server/src/main/resources/application.yaml +++ b/server/src/main/resources/application.yaml @@ -72,6 +72,13 @@ 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. For now it is only sent as + # the share_identifier job parameter, since the queue has no fair-share policy yet. + share: + small-max-size: 50MB + small: small + large: large wfs-default-param: fields: From 84fcfe313f508478092846fc48fe4acdb0274677 Mon Sep 17 00:00:00 2001 From: Yuxuan HU Date: Wed, 7 Oct 2026 15:37:18 +1100 Subject: [PATCH 3/7] add shared share_identifier to job parameter --- .../processes/DownloadAdmissionService.java | 37 +++++++++++--- .../DownloadAdmissionServiceTest.java | 49 ++++++++++++++++++- 2 files changed, 78 insertions(+), 8 deletions(-) diff --git a/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionService.java b/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionService.java index f65eeafa..e55f59cd 100644 --- a/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionService.java +++ b/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionService.java @@ -2,6 +2,7 @@ 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; @@ -9,6 +10,7 @@ import java.util.HashMap; import java.util.Map; +import java.util.OptionalLong; import java.util.concurrent.locks.ReentrantLock; /** @@ -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 @@ -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(); @@ -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; } /** @@ -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. @@ -73,6 +77,10 @@ public String submit(DownloadRequest request) throws JsonProcessingException { try { Map parameters = restServices.buildDownloadParameters(request); + String shareIdentifier = shareFor(estimatedBytes); + log.info("Download for uuid {} estimated {} bytes, share {}", request.uuid(), + estimatedBytes.isPresent() ? estimatedBytes.getAsLong() : "none", shareIdentifier); + parameters.put(DatasetDownloadEnums.Parameter.SHARE_IDENTIFIER.getValue(), shareIdentifier); String jobName = RestServices.downloadJobName(request.recipient()); String awsJobId = restServices.submitDownloadJob(jobName, parameters); counter.recordSubmitted(awsJobId, request.recipient()); @@ -85,7 +93,10 @@ 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); @@ -93,7 +104,7 @@ private void rejectIfTooLarge(DownloadRequest request) { // 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(); @@ -102,6 +113,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) { diff --git a/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java b/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java index cbdb1135..6a0ae509 100644 --- a/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java +++ b/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java @@ -6,6 +6,7 @@ import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; import org.springframework.util.unit.DataSize; @@ -63,7 +64,16 @@ private DownloadAdmissionService build(DownloadLimitProperties limits) { } private DownloadAdmissionService build(DownloadLimitProperties limits, DownloadSizeLimitProperties sizeLimit) { - return new DownloadAdmissionService(restServices, counter, limits, sizeLimit); + return new DownloadAdmissionService(restServices, counter, limits, sizeLimit, + new DownloadShareProperties(DataSize.ofMegabytes(50), "small", "large")); + } + + /** The share_identifier parameter of the one download submitted to Batch. */ + private String submittedShare() { + @SuppressWarnings("unchecked") + ArgumentCaptor> parameters = ArgumentCaptor.forClass(Map.class); + verify(restServices).submitDownloadJob(anyString(), parameters.capture()); + return parameters.getValue().get("share_identifier"); } private static DownloadLimitProperties limits(boolean enabled, int maxConcurrent) { @@ -232,4 +242,41 @@ void theSizeLimitCanBeTurnedOff() throws Exception { assertNotNull(disabled.submit(request(RECIPIENT))); verify(restServices, never()).estimateDownloadBytes(any()); } + + @Test + void aDownloadUnderTheShareThresholdIsTaggedSmall() throws Exception { + when(restServices.estimateDownloadBytes(any())).thenReturn(DataSize.ofMegabytes(50).toBytes() - 1); + + service.submit(request(RECIPIENT)); + + assertEquals("small", submittedShare()); + } + + @Test + void aDownloadAtTheShareThresholdIsTaggedLarge() throws Exception { + when(restServices.estimateDownloadBytes(any())).thenReturn(DataSize.ofMegabytes(50).toBytes()); + + service.submit(request(RECIPIENT)); + + assertEquals("large", submittedShare()); + } + + @Test + void aDownloadWithAFailedEstimateIsTaggedLarge() throws Exception { + when(restServices.estimateDownloadBytes(any())).thenThrow(new RuntimeException("DAS is down")); + + service.submit(request(RECIPIENT)); + + assertEquals("large", submittedShare()); + } + + @Test + void aDownloadIsTaggedLargeWhenTheSizeLimitIsOff() throws Exception { + DownloadAdmissionService disabled = build(limits(true, 10), sizeLimit(false)); + + disabled.submit(request(RECIPIENT)); + + assertEquals("large", submittedShare()); + verify(restServices, never()).estimateDownloadBytes(any()); + } } From 0837043cc142a37f1f05848c69ffdc9133605307 Mon Sep 17 00:00:00 2001 From: Yuxuan HU Date: Fri, 9 Oct 2026 11:28:29 +1100 Subject: [PATCH 4/7] add shareIdentifier for submitJob --- .../server/processes/DownloadAdmissionService.java | 4 +++- .../server/processes/DownloadShareProperties.java | 8 ++++++-- .../org/aodn/ogcapi/server/processes/RestServices.java | 10 +++++++--- server/src/main/resources/application.yaml | 8 ++++++-- 4 files changed, 22 insertions(+), 8 deletions(-) diff --git a/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionService.java b/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionService.java index e55f59cd..e6bb4339 100644 --- a/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionService.java +++ b/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionService.java @@ -80,9 +80,11 @@ public String submit(DownloadRequest request) throws JsonProcessingException { 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; diff --git a/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadShareProperties.java b/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadShareProperties.java index c744bb42..0e4e656f 100644 --- a/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadShareProperties.java +++ b/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadShareProperties.java @@ -8,11 +8,15 @@ * 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 values go * into the Batch shareIdentifier later, which allows only alphanumerics. + * When enabled is true, the share is also sent as the Batch shareIdentifier. 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 + @DefaultValue("small-downloads") String small, + @DefaultValue("large-downloads") String large ) { } diff --git a/server/src/main/java/au/org/aodn/ogcapi/server/processes/RestServices.java b/server/src/main/java/au/org/aodn/ogcapi/server/processes/RestServices.java index b8cababf..90f1543f 100644 --- a/server/src/main/java/au/org/aodn/ogcapi/server/processes/RestServices.java +++ b/server/src/main/java/au/org/aodn/ogcapi/server/processes/RestServices.java @@ -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 parameters) { - String jobId = submitJob(jobName, this.batchJobQueue, this.batchJobDefinition, parameters); + public String submitDownloadJob(String jobName, Map 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 parameters) { + private String submitJob(String jobName, String jobQueue, String jobDefinition, Map 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 @@ -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); diff --git a/server/src/main/resources/application.yaml b/server/src/main/resources/application.yaml index a577b74b..f0669965 100644 --- a/server/src/main/resources/application.yaml +++ b/server/src/main/resources/application.yaml @@ -73,9 +73,13 @@ aws: 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. For now it is only sent as - # the share_identifier job parameter, since the queue has no fair-share policy yet. + # 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 From 435058dde607879306f31dee6efa2f336a74564f Mon Sep 17 00:00:00 2001 From: Yuxuan HU Date: Fri, 9 Oct 2026 11:36:21 +1100 Subject: [PATCH 5/7] add tests --- .../DownloadAdmissionServiceTest.java | 62 +++++++++++++++---- .../server/processes/RestServicesTest.java | 24 ++++++- 2 files changed, 74 insertions(+), 12 deletions(-) diff --git a/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java b/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java index 6a0ae509..f6781af7 100644 --- a/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java +++ b/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java @@ -19,6 +19,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; @@ -51,7 +52,7 @@ class DownloadAdmissionServiceTest { void setUp() throws JsonProcessingException { lenient().when(restServices.buildDownloadParameters(any())) .thenAnswer(invocation -> new HashMap()); - lenient().when(restServices.submitDownloadJob(anyString(), any())) + lenient().when(restServices.submitDownloadJob(anyString(), any(), any())) .thenAnswer(invocation -> UUID.randomUUID().toString()); lenient().when(counter.countInFlight(anyString())) .thenAnswer(invocation -> current(invocation.getArgument(0)).get()); @@ -64,18 +65,30 @@ private DownloadAdmissionService build(DownloadLimitProperties limits) { } private DownloadAdmissionService build(DownloadLimitProperties limits, DownloadSizeLimitProperties sizeLimit) { + return build(limits, sizeLimit, false); + } + + private DownloadAdmissionService build(DownloadLimitProperties limits, DownloadSizeLimitProperties sizeLimit, + boolean shareEnabled) { return new DownloadAdmissionService(restServices, counter, limits, sizeLimit, - new DownloadShareProperties(DataSize.ofMegabytes(50), "small", "large")); + new DownloadShareProperties(shareEnabled, DataSize.ofMegabytes(50), "small-downloads", "large-downloads")); } /** The share_identifier parameter of the one download submitted to Batch. */ private String submittedShare() { @SuppressWarnings("unchecked") ArgumentCaptor> parameters = ArgumentCaptor.forClass(Map.class); - verify(restServices).submitDownloadJob(anyString(), parameters.capture()); + verify(restServices).submitDownloadJob(anyString(), parameters.capture(), any()); return parameters.getValue().get("share_identifier"); } + /** The top-level shareIdentifier argument of the one download submitted to Batch. */ + private String submittedShareIdentifier() { + ArgumentCaptor shareIdentifier = ArgumentCaptor.forClass(String.class); + verify(restServices).submitDownloadJob(anyString(), any(), shareIdentifier.capture()); + return shareIdentifier.getValue(); + } + private static DownloadLimitProperties limits(boolean enabled, int maxConcurrent) { return new DownloadLimitProperties(enabled, maxConcurrent, Duration.ofSeconds(15)); } @@ -102,7 +115,7 @@ void underTheLimitSubmitsAndReturnsTheAwsJobId() throws Exception { String jobId = service.submit(request(RECIPIENT)); assertNotNull(jobId); - verify(restServices).submitDownloadJob(eq(RECIPIENT_JOB_NAME), any()); + verify(restServices).submitDownloadJob(eq(RECIPIENT_JOB_NAME), any(), any()); verify(restServices).notifyUser(any(), any(), any(), any(), any(), any(), any(), any(), any(), any()); } @@ -112,7 +125,7 @@ void atTheLimitIsRejectedOutright() throws Exception { assertThrows(DownloadLimitExceededException.class, () -> service.submit(request(RECIPIENT))); - verify(restServices, never()).submitDownloadJob(anyString(), any()); + verify(restServices, never()).submitDownloadJob(anyString(), any(), any()); // Nothing was submitted, so the user must not be told their file is being produced. verify(restServices, never()).notifyUser(any(), any(), any(), any(), any(), any(), any(), any(), any(), any()); } @@ -169,7 +182,7 @@ void theLimitCanBeTurnedOffEntirely() throws Exception { String jobId = disabled.submit(request(RECIPIENT)); assertNotNull(jobId); - verify(restServices).submitDownloadJob(anyString(), any()); + verify(restServices).submitDownloadJob(anyString(), any(), any()); verify(counter, never()).refreshIfStale(); } @@ -178,7 +191,7 @@ void aFailedSubmitDoesNotLeakAReservedSlot() throws Exception { // One slot free; the counter will not move because the submit below never actually // reaches AWS. current(RECIPIENT).set(9); - org.mockito.Mockito.when(restServices.submitDownloadJob(anyString(), any())) + org.mockito.Mockito.when(restServices.submitDownloadJob(anyString(), any(), any())) .thenThrow(new IllegalStateException("AWS Batch rejected the job")); assertThrows(IllegalStateException.class, () -> service.submit(request(RECIPIENT))); @@ -188,7 +201,7 @@ void aFailedSubmitDoesNotLeakAReservedSlot() throws Exception { // failed attempt never actually started a download. org.mockito.Mockito.reset(restServices); lenient().when(restServices.buildDownloadParameters(any())).thenReturn(new HashMap<>()); - lenient().when(restServices.submitDownloadJob(anyString(), any())) + lenient().when(restServices.submitDownloadJob(anyString(), any(), any())) .thenReturn(UUID.randomUUID().toString()); assertNotNull(service.submit(request(RECIPIENT))); @@ -205,7 +218,7 @@ void aDownloadAtTheSizeLimitIsRejectedBeforeSubmitAndEmail() throws Exception { assertEquals("Download is unavailable because the selected dataset is too large " + "(estimated 180 GB, limit 180 GB). Please refine your selection to reduce the dataset size.", exception.getMessage()); - verify(restServices, never()).submitDownloadJob(anyString(), any()); + verify(restServices, never()).submitDownloadJob(anyString(), any(), any()); verify(restServices, never()).notifyUser(any(), any(), any(), any(), any(), any(), any(), any(), any(), any()); } @@ -214,7 +227,7 @@ void aDownloadJustUnderTheSizeLimitIsSubmitted() throws Exception { when(restServices.estimateDownloadBytes(any())).thenReturn(MAX_BYTES - 1); assertNotNull(service.submit(request(RECIPIENT))); - verify(restServices).submitDownloadJob(eq(RECIPIENT_JOB_NAME), any()); + verify(restServices).submitDownloadJob(eq(RECIPIENT_JOB_NAME), any(), any()); } @Test @@ -232,7 +245,7 @@ void aFailedEstimateStillSubmits() throws Exception { when(restServices.estimateDownloadBytes(any())).thenThrow(new RuntimeException("DAS is down")); assertNotNull(service.submit(request(RECIPIENT))); - verify(restServices).submitDownloadJob(eq(RECIPIENT_JOB_NAME), any()); + verify(restServices).submitDownloadJob(eq(RECIPIENT_JOB_NAME), any(), any()); } @Test @@ -279,4 +292,31 @@ void aDownloadIsTaggedLargeWhenTheSizeLimitIsOff() throws Exception { assertEquals("large", submittedShare()); verify(restServices, never()).estimateDownloadBytes(any()); } + + @Test + void whenShareIsEnabledTheShareIdentifierMatchesTheParameter() throws Exception { + // Under the threshold tags small; with the switch on, that share is also sent as the + // top-level Batch shareIdentifier. + DownloadAdmissionService shareEnabled = build(limits(true, 10), sizeLimit(true), true); + when(restServices.estimateDownloadBytes(any())) + .thenReturn(DataSize.ofMegabytes(50).toBytes() - 1); + + shareEnabled.submit(request(RECIPIENT)); + + assertEquals("small-downloads", submittedShare()); + assertEquals("small-downloads", submittedShareIdentifier()); + } + + @Test + void whenShareIsDisabledNoShareIdentifierIsSentButTheParameterRemains() throws Exception { + // Default build has the switch off: the share_identifier parameter still goes out for + // DAS, but no top-level shareIdentifier is sent so a FIFO queue does not reject it. + when(restServices.estimateDownloadBytes(any())) + .thenReturn(DataSize.ofMegabytes(50).toBytes() - 1); + + service.submit(request(RECIPIENT)); + + assertEquals("small-downloads", submittedShare()); + assertNull(submittedShareIdentifier()); + } } diff --git a/server/src/test/java/au/org/aodn/ogcapi/server/processes/RestServicesTest.java b/server/src/test/java/au/org/aodn/ogcapi/server/processes/RestServicesTest.java index 7512bf81..51519496 100644 --- a/server/src/test/java/au/org/aodn/ogcapi/server/processes/RestServicesTest.java +++ b/server/src/test/java/au/org/aodn/ogcapi/server/processes/RestServicesTest.java @@ -16,7 +16,10 @@ import software.amazon.awssdk.services.batch.model.SubmitJobRequest; import software.amazon.awssdk.services.batch.model.SubmitJobResponse; +import java.util.HashMap; + import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyMap; @@ -54,7 +57,7 @@ private String downloadData(String uuid, String key, String startDate, String en DownloadRequest request = new DownloadRequest(uuid, key, startDate, endDate, polygons, recipient, collectionTitle, fullMetadataLink, suggestedCitation, outputFormat); return restServices.submitDownloadJob( - RestServices.downloadJobName(recipient), restServices.buildDownloadParameters(request)); + RestServices.downloadJobName(recipient), restServices.buildDownloadParameters(request), null); } @Test @@ -112,9 +115,28 @@ public void testDownloadDataCapturesSubmitJobRequest() throws JsonProcessingExce assertEquals("geotiff", captured.parameters().get("output_format")); assertEquals("https://metadata.imas.utas.edu.au/.../test-uuid-123", captured.parameters().get("full_metadata_link")); + // The helper passes null as the share, so no top-level shareIdentifier goes out. + assertNull(captured.shareIdentifier()); assertEquals(jobId, response); } + @Test + public void submitDownloadJobSetsTheShareIdentifierOnTheRequest() throws JsonProcessingException { + // Arrange + String jobId = "12345"; + SubmitJobResponse submitJobResponse = SubmitJobResponse.builder().jobId(jobId).build(); + when(batchClient.submitJob(any(SubmitJobRequest.class))).thenReturn(submitJobResponse); + + // Act: submit with a non-null share, as the admission service does on a fair-share queue. + restServices.submitDownloadJob("test-job", new HashMap<>(), "small-downloads"); + + // Capture the submitted request + ArgumentCaptor captor = ArgumentCaptor.forClass(SubmitJobRequest.class); + verify(batchClient, times(1)).submitJob(captor.capture()); + + assertEquals("small-downloads", captor.getValue().shareIdentifier()); + } + @Test public void submitJobReplacesEmptySuggestedCitationWithUnavailable() throws JsonProcessingException { // Arrange From d1890f381aba42723b64ce54a97e835b82a3b5b2 Mon Sep 17 00:00:00 2001 From: Yuxuan HU Date: Fri, 9 Oct 2026 11:55:10 +1100 Subject: [PATCH 6/7] rename shareIdentifier labels --- .../server/processes/DownloadShareProperties.java | 14 +++++++------- .../processes/DownloadAdmissionServiceTest.java | 2 +- .../ogcapi/server/processes/RestServicesTest.java | 4 ++-- 3 files changed, 10 insertions(+), 10 deletions(-) diff --git a/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadShareProperties.java b/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadShareProperties.java index 0e4e656f..4c1a9888 100644 --- a/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadShareProperties.java +++ b/server/src/main/java/au/org/aodn/ogcapi/server/processes/DownloadShareProperties.java @@ -6,17 +6,17 @@ /** * 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 values go - * into the Batch shareIdentifier later, which allows only alphanumerics. - * When enabled is true, the share is also sent as the Batch shareIdentifier. 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. + * 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-downloads") String small, - @DefaultValue("large-downloads") String large + @DefaultValue("small") String small, + @DefaultValue("large") String large ) { } diff --git a/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java b/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java index f6781af7..bec67e45 100644 --- a/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java +++ b/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java @@ -71,7 +71,7 @@ private DownloadAdmissionService build(DownloadLimitProperties limits, DownloadS private DownloadAdmissionService build(DownloadLimitProperties limits, DownloadSizeLimitProperties sizeLimit, boolean shareEnabled) { return new DownloadAdmissionService(restServices, counter, limits, sizeLimit, - new DownloadShareProperties(shareEnabled, DataSize.ofMegabytes(50), "small-downloads", "large-downloads")); + new DownloadShareProperties(shareEnabled, DataSize.ofMegabytes(50), "small", "large")); } /** The share_identifier parameter of the one download submitted to Batch. */ diff --git a/server/src/test/java/au/org/aodn/ogcapi/server/processes/RestServicesTest.java b/server/src/test/java/au/org/aodn/ogcapi/server/processes/RestServicesTest.java index 51519496..1ffedd25 100644 --- a/server/src/test/java/au/org/aodn/ogcapi/server/processes/RestServicesTest.java +++ b/server/src/test/java/au/org/aodn/ogcapi/server/processes/RestServicesTest.java @@ -128,13 +128,13 @@ public void submitDownloadJobSetsTheShareIdentifierOnTheRequest() throws JsonPro when(batchClient.submitJob(any(SubmitJobRequest.class))).thenReturn(submitJobResponse); // Act: submit with a non-null share, as the admission service does on a fair-share queue. - restServices.submitDownloadJob("test-job", new HashMap<>(), "small-downloads"); + restServices.submitDownloadJob("test-job", new HashMap<>(), "small"); // Capture the submitted request ArgumentCaptor captor = ArgumentCaptor.forClass(SubmitJobRequest.class); verify(batchClient, times(1)).submitJob(captor.capture()); - assertEquals("small-downloads", captor.getValue().shareIdentifier()); + assertEquals("small", captor.getValue().shareIdentifier()); } @Test From f2171aae5bb3e92a9d4133638e2542d4ca1f091d Mon Sep 17 00:00:00 2001 From: Yuxuan HU Date: Fri, 9 Oct 2026 11:56:16 +1100 Subject: [PATCH 7/7] rename shareIdentifier labels --- .../server/processes/DownloadAdmissionServiceTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java b/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java index bec67e45..a9408bfd 100644 --- a/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java +++ b/server/src/test/java/au/org/aodn/ogcapi/server/processes/DownloadAdmissionServiceTest.java @@ -303,8 +303,8 @@ void whenShareIsEnabledTheShareIdentifierMatchesTheParameter() throws Exception shareEnabled.submit(request(RECIPIENT)); - assertEquals("small-downloads", submittedShare()); - assertEquals("small-downloads", submittedShareIdentifier()); + assertEquals("small", submittedShare()); + assertEquals("small", submittedShareIdentifier()); } @Test @@ -316,7 +316,7 @@ void whenShareIsDisabledNoShareIdentifierIsSentButTheParameterRemains() throws E service.submit(request(RECIPIENT)); - assertEquals("small-downloads", submittedShare()); + assertEquals("small", submittedShare()); assertNull(submittedShareIdentifier()); } }