Skip to content

Commit 2040eb1

Browse files
committed
fix: honor queue-level DelaySeconds for standard SQS queues
1 parent 976f819 commit 2040eb1

4 files changed

Lines changed: 76 additions & 19 deletions

File tree

src/main/java/io/github/hectorvent/floci/services/sqs/SqsJsonHandler.java

Lines changed: 16 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -195,7 +195,8 @@ private Response handleGetQueueAttributes(JsonNode request, String region) {
195195
private Response handleSendMessage(JsonNode request, String region) {
196196
String queueUrl = request.path("QueueUrl").asText(null);
197197
String messageBody = request.path("MessageBody").asText(null);
198-
int delaySeconds = request.path("DelaySeconds").asInt(0);
198+
JsonNode delayNode = request.path("DelaySeconds");
199+
Integer delaySeconds = parseOptionalInteger(delayNode, "DelaySeconds");
199200
String messageGroupId = request.path("MessageGroupId").asText(null);
200201
String messageDeduplicationId = request.path("MessageDeduplicationId").asText(null);
201202

@@ -347,7 +348,7 @@ private Response handleSendMessageBatch(JsonNode request, String region) {
347348
ArrayNode successful = objectMapper.createArrayNode();
348349
ArrayNode failed = objectMapper.createArrayNode();
349350

350-
record ParsedEntry(String id, String body, int delay, String groupId, String dedupId,
351+
record ParsedEntry(String id, String body, Integer delay, String groupId, String dedupId,
351352
Map<String, MessageAttributeValue> attributes, String awsTraceHeader) {}
352353

353354
List<ParsedEntry> parsedEntries = new ArrayList<>();
@@ -356,7 +357,8 @@ record ParsedEntry(String id, String body, int delay, String groupId, String ded
356357
for (JsonNode entry : entries) {
357358
String id = entry.path("Id").asText();
358359
String messageBody = entry.path("MessageBody").asText(null);
359-
int delaySeconds = entry.path("DelaySeconds").asInt(0);
360+
JsonNode entryDelayNode = entry.path("DelaySeconds");
361+
Integer delaySeconds = parseOptionalInteger(entryDelayNode, "DelaySeconds");
360362
String messageGroupId = entry.path("MessageGroupId").asText(null);
361363
String messageDeduplicationId = entry.path("MessageDeduplicationId").asText(null);
362364

@@ -583,4 +585,15 @@ private Map<String, String> jsonNodeToMap(JsonNode node) {
583585
}
584586
return map;
585587
}
588+
589+
private Integer parseOptionalInteger(JsonNode node, String paramName) {
590+
if (node.isMissingNode() || node.isNull()) {
591+
return null;
592+
}
593+
if (!node.isIntegralNumber()) {
594+
throw new AwsException("InvalidParameterValue",
595+
"Value for parameter " + paramName + " is invalid. Reason: Must be an integer.", 400);
596+
}
597+
return node.asInt();
598+
}
586599
}

src/main/java/io/github/hectorvent/floci/services/sqs/SqsQueryHandler.java

Lines changed: 14 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -196,7 +196,7 @@ private Response handleGetQueueAttributes(MultivaluedMap<String, String> params,
196196
private Response handleSendMessage(MultivaluedMap<String, String> params, String region) {
197197
String queueUrl = getParam(params, "QueueUrl");
198198
String body = getParam(params, "MessageBody");
199-
int delaySeconds = getIntParam(params, "DelaySeconds", 0);
199+
Integer delaySeconds = getIntegerParam(params, "DelaySeconds");
200200
String messageGroupId = getParam(params, "MessageGroupId");
201201
String messageDeduplicationId = getParam(params, "MessageDeduplicationId");
202202

@@ -331,7 +331,7 @@ private Response handleSendMessageBatch(MultivaluedMap<String, String> params, S
331331
String queueUrl = getParam(params, "QueueUrl");
332332
var xml = new XmlBuilder();
333333

334-
record ParsedEntry(String id, String body, int delay, String groupId, String dedupId,
334+
record ParsedEntry(String id, String body, Integer delay, String groupId, String dedupId,
335335
Map<String, MessageAttributeValue> attributes, String awsTraceHeader) {}
336336

337337
List<ParsedEntry> parsedEntries = new ArrayList<>();
@@ -340,7 +340,7 @@ record ParsedEntry(String id, String body, int delay, String groupId, String ded
340340
String id = getParam(params, "SendMessageBatchRequestEntry." + i + ".Id");
341341
if (id == null) break;
342342
String body = getParam(params, "SendMessageBatchRequestEntry." + i + ".MessageBody");
343-
int delaySeconds = getIntParam(params, "SendMessageBatchRequestEntry." + i + ".DelaySeconds", 0);
343+
Integer delaySeconds = getIntegerParam(params, "SendMessageBatchRequestEntry." + i + ".DelaySeconds");
344344
String messageGroupId = getParam(params, "SendMessageBatchRequestEntry." + i + ".MessageGroupId");
345345
String messageDeduplicationId = getParam(params, "SendMessageBatchRequestEntry." + i + ".MessageDeduplicationId");
346346

@@ -560,6 +560,17 @@ private int getIntParam(MultivaluedMap<String, String> params, String name, int
560560
}
561561
}
562562

563+
private Integer getIntegerParam(MultivaluedMap<String, String> params, String name) {
564+
String value = params.getFirst(name);
565+
if (value == null) return null;
566+
try {
567+
return Integer.parseInt(value);
568+
} catch (NumberFormatException e) {
569+
throw new AwsException("InvalidParameterValue",
570+
"Value for parameter " + name + " is invalid. Reason: Must be an integer.", 400);
571+
}
572+
}
573+
563574
private Map<String, String> extractAttributes(MultivaluedMap<String, String> params) {
564575
Map<String, String> attributes = new HashMap<>();
565576
for (int i = 1; ; i++) {

src/main/java/io/github/hectorvent/floci/services/sqs/SqsService.java

Lines changed: 14 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -340,25 +340,25 @@ public Map<String, String> getQueueAttributes(String queueUrl, List<String> attr
340340
return filtered;
341341
}
342342

343-
public Message sendMessage(String queueUrl, String body, int delaySeconds, String region) {
343+
public Message sendMessage(String queueUrl, String body, Integer delaySeconds, String region) {
344344
return sendMessage(queueUrl, body, delaySeconds, null, null, region);
345345
}
346346

347-
public Message sendMessage(String queueUrl, String body, int delaySeconds,
347+
public Message sendMessage(String queueUrl, String body, Integer delaySeconds,
348348
String messageGroupId, String messageDeduplicationId,
349349
String region) {
350350
return sendMessage(queueUrl, body, delaySeconds, messageGroupId, messageDeduplicationId, null, region);
351351
}
352352

353-
public Message sendMessage(String queueUrl, String body, int delaySeconds,
353+
public Message sendMessage(String queueUrl, String body, Integer delaySeconds,
354354
String messageGroupId, String messageDeduplicationId,
355355
Map<String, MessageAttributeValue> messageAttributes,
356356
String region) {
357357
return sendMessage(queueUrl, body, delaySeconds, messageGroupId, messageDeduplicationId,
358358
messageAttributes, null, region);
359359
}
360360

361-
public Message sendMessage(String queueUrl, String body, int delaySeconds,
361+
public Message sendMessage(String queueUrl, String body, Integer delaySeconds,
362362
String messageGroupId, String messageDeduplicationId,
363363
Map<String, MessageAttributeValue> messageAttributes,
364364
String awsTraceHeader,
@@ -381,15 +381,16 @@ public Message sendMessage(String queueUrl, String body, int delaySeconds,
381381
// Resolve the effective delay:
382382
// - FIFO queues only support queue-level DelaySeconds per AWS SQS,
383383
// so any per-message value is ignored and we always use the
384-
// queue attribute. Without this, FIFO silently dropped the
385-
// queue-level default (issue #475).
386-
// - Standard queues honor per-message DelaySeconds when provided
387-
// (> 0). Applying the queue-level default on the standard path
388-
// requires distinguishing "omitted" from "explicit 0" in the
389-
// handlers, which the current int-parameter API cannot express;
390-
// that's left as follow-up work -- this patch only addresses
391-
// the FIFO regression called out in the issue.
392-
int effectiveDelaySeconds = queue.isFifo() ? queueDelaySeconds : delaySeconds;
384+
// queue attribute.
385+
// - Standard queues honor per-message DelaySeconds when explicitly
386+
// provided (non-null). When omitted (null), the queue-level
387+
// default applies.
388+
int effectiveDelaySeconds;
389+
if (queue.isFifo()) {
390+
effectiveDelaySeconds = queueDelaySeconds;
391+
} else {
392+
effectiveDelaySeconds = (delaySeconds != null) ? delaySeconds : queueDelaySeconds;
393+
}
393394

394395
// FIFO queue validation
395396
if (queue.isFifo()) {

src/test/java/io/github/hectorvent/floci/services/sqs/SqsServiceTest.java

Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -461,6 +461,38 @@ void fifoQueueIgnoresPerMessageDelaySeconds() {
461461
"FIFO queues must ignore per-message DelaySeconds");
462462
}
463463

464+
// --- Queue-level DelaySeconds for standard queues ---
465+
466+
@Test
467+
void queueLevelDelaySecondsAppliesToStandardQueue() {
468+
String region = "eu-west-1";
469+
Queue queue = sqsService.createQueue("delay-standard",
470+
Map.of("DelaySeconds", "1"), region);
471+
sqsService.sendMessage(queue.getQueueUrl(), "msg", null, region);
472+
473+
List<Message> immediate = sqsService.receiveMessage(queue.getQueueUrl(), 1, 0, 0, region);
474+
assertTrue(immediate.isEmpty(),
475+
"Standard queue should honor queue-level DelaySeconds");
476+
477+
try { Thread.sleep(1100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); }
478+
479+
List<Message> later = sqsService.receiveMessage(queue.getQueueUrl(), 1, 0, 0, region);
480+
assertEquals(1, later.size(),
481+
"Message should become visible once DelaySeconds elapses");
482+
}
483+
484+
@Test
485+
void explicitZeroDelayOverridesQueueDefault() {
486+
String region = "eu-west-1";
487+
Queue queue = sqsService.createQueue("delay-override",
488+
Map.of("DelaySeconds", "10"), region);
489+
sqsService.sendMessage(queue.getQueueUrl(), "msg", 0, region);
490+
491+
List<Message> immediate = sqsService.receiveMessage(queue.getQueueUrl(), 1, 0, 0, region);
492+
assertEquals(1, immediate.size(),
493+
"Explicit DelaySeconds=0 must override queue-level default");
494+
}
495+
464496
// --- clearFifoDeduplicationCacheOnPurge tests ---
465497

466498
@Test

0 commit comments

Comments
 (0)