#154: pin AMQP reply inbox prefetch #181
@@ -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;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user