From fa3f910d44c3d9551199bd937fc8570bc968c3a1 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Fri, 28 Aug 2026 06:03:54 +0700 Subject: [PATCH] #154: pin AMQP reply inbox prefetch --- .../dev/ltms/fleet/msg/AmqpReplyInbox.java | 5 +- .../fleet/msg/AmqpReplyInboxPrefetchTest.java | 149 ++++++++++++++++++ 2 files changed, 152 insertions(+), 2 deletions(-) create mode 100644 fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxPrefetchTest.java diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java b/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java index bbb2d9f..dfef5c2 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/AmqpReplyInbox.java @@ -29,8 +29,9 @@ import java.util.concurrent.TimeoutException; * *

Mapping — consume-and-hold with deferred manual ack. Each target has a durable * queue {@code agent..inbox}. The gateway that owns the target starts a manual-ack consumer - * ({@link #own}) that pulls persistent messages off that queue into an in-memory held map - * (keyed by {@code msgId}) but does not ack them. {@link #peek} returns that snapshot; + * ({@link #own}) that pulls persistent messages, up to its prefetch window, off that queue into an + * in-memory held map (keyed by {@code msgId}) but does not ack them. + * {@link #peek} returns that snapshot; * {@link #ack} acks the broker delivery-tag and drops the entry. Because messages stay unacked until * the owning gateway actually drains them, a crash (or a {@code java -jar} bounce) before caller-ack * leaves them on the broker — it redelivers on reconnect. That is the durability the in-memory diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxPrefetchTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxPrefetchTest.java new file mode 100644 index 0000000..b5dfb8f --- /dev/null +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/AmqpReplyInboxPrefetchTest.java @@ -0,0 +1,149 @@ +package dev.ltms.fleet.msg; + +import com.rabbitmq.client.AMQP; +import com.rabbitmq.client.Channel; +import com.rabbitmq.client.Connection; +import com.rabbitmq.client.DeliverCallback; +import com.rabbitmq.client.Delivery; +import com.rabbitmq.client.Envelope; +import org.junit.jupiter.api.Test; + +import java.io.IOException; +import java.lang.reflect.InvocationHandler; +import java.lang.reflect.Proxy; +import java.nio.charset.StandardCharsets; +import java.util.ArrayDeque; +import java.util.LinkedHashMap; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; + +/** + * Pins the manual-ack prefetch behaviour without a broker. The fake channel models a broker that + * sends no more than its QoS window of unacked deliveries. If {@link AmqpReplyInbox} starts acking + * messages while it adds them to {@code held}, this test drains the whole fake queue instead. + */ +class AmqpReplyInboxPrefetchTest { + + @Test + void unackedDeliveriesKeepTheHeldBacklogAtThePrefetchWindow() { + int prefetch = 3; + int published = 8; + PrefetchBroker broker = new PrefetchBroker(); + for (int i = 0; i < published; i++) { + broker.publish("m" + i, "payload " + i); + } + + try (AmqpReplyInbox inbox = new AmqpReplyInbox(connectionFor(broker.channel()), prefetch)) { + inbox.own("worker"); + + assertEquals(prefetch, inbox.peek("worker").size(), + "held messages must stop at the unacked prefetch window"); + assertEquals(published - prefetch, broker.queuedCount(), + "messages beyond the window must remain on the broker"); + assertEquals(0, broker.ackCount(), "receipt must not ack a held message"); + + inbox.ack("worker", "m0"); + + assertEquals(prefetch, inbox.peek("worker").size(), + "one caller ack frees exactly one slot for the broker"); + assertEquals(published - prefetch - 1, broker.queuedCount(), + "only one queued message may enter after one caller ack"); + assertEquals(1, broker.ackCount(), "only the caller ack may reach the broker"); + } + } + + private static Connection connectionFor(Channel consumeChannel) { + Channel publishChannel = (Channel) Proxy.newProxyInstance( + AmqpReplyInboxPrefetchTest.class.getClassLoader(), new Class[] {Channel.class}, + (proxy, method, args) -> defaultValue(method.getReturnType())); + AtomicInteger channelCalls = new AtomicInteger(); + InvocationHandler handler = (proxy, method, args) -> { + if (method.getName().equals("createChannel") && (args == null || args.length == 0)) { + return channelCalls.getAndIncrement() == 0 ? consumeChannel : publishChannel; + } + return defaultValue(method.getReturnType()); + }; + return (Connection) Proxy.newProxyInstance(AmqpReplyInboxPrefetchTest.class.getClassLoader(), + new Class[] {Connection.class}, handler); + } + + private static final class PrefetchBroker implements InvocationHandler { + private final ArrayDeque queued = new ArrayDeque<>(); + private final Map unacked = new LinkedHashMap<>(); + private DeliverCallback consumer; + private int prefetch; + private int acks; + private long nextTag = 1; + + Channel channel() { + return (Channel) Proxy.newProxyInstance(AmqpReplyInboxPrefetchTest.class.getClassLoader(), + new Class[] {Channel.class}, this); + } + + void publish(String msgId, String content) { + queued.add(new Delivery(new Envelope(nextTag++, false, "", ""), + new AMQP.BasicProperties.Builder().messageId(msgId).build(), + content.getBytes(StandardCharsets.UTF_8))); + } + + int queuedCount() { + return queued.size(); + } + + int ackCount() { + return acks; + } + + @Override + public Object invoke(Object proxy, java.lang.reflect.Method method, Object[] args) throws IOException { + switch (method.getName()) { + case "basicQos" -> { + prefetch = (int) args[0]; + return null; + } + case "basicConsume" -> { + assertFalse((boolean) args[1], "the inbox consumer must use manual acknowledgements"); + consumer = (DeliverCallback) args[2]; + deliverAvailable(); + return "consumer"; + } + case "basicAck" -> { + unacked.remove((long) args[0]); + acks++; + deliverAvailable(); + return null; + } + default -> { + return defaultValue(method.getReturnType()); + } + } + } + + private void deliverAvailable() throws IOException { + while (consumer != null && unacked.size() < prefetch && !queued.isEmpty()) { + Delivery delivery = queued.removeFirst(); + unacked.put(delivery.getEnvelope().getDeliveryTag(), delivery); + consumer.handle("consumer", delivery); + } + } + } + + private static Object defaultValue(Class type) { + if (!type.isPrimitive() || type == void.class) { + return null; + } + if (type == boolean.class) { + return false; + } + if (type == long.class) { + return 0L; + } + if (type == int.class) { + return 0; + } + return 0; + } +} -- 2.52.0