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;
+ }
+}