#154: pin AMQP reply inbox prefetch #181

Merged
ltms merged 1 commits from worker/fleetd-154-held-bound-c0f2e4-7 into main 2026-08-28 01:06:43 +02:00
2 changed files with 152 additions and 2 deletions
@@ -29,8 +29,9 @@ import java.util.concurrent.TimeoutException;
*
* <p><strong>Mapping — consume-and-hold with deferred manual ack.</strong> Each target has a durable
* queue {@code agent.<target>.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 <em>held</em> map
* (keyed by {@code msgId}) but does <em>not</em> 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 <em>held</em> map (keyed by {@code msgId}) but does <em>not</em> 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
@@ -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<Delivery> queued = new ArrayDeque<>();
private final Map<Long, Delivery> 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;
}
}