Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 7b98cca967 |
@@ -654,6 +654,23 @@ guard:
|
||||
# uriEnv: LAVINMQ_URI
|
||||
# prefetch: 32
|
||||
|
||||
# Shared cross-host LEADER coordination broker. OMIT this block to leave lead-to-lead messaging
|
||||
# off entirely (config-only in this ticket — nothing here wires it into a live LeadMailbox yet).
|
||||
# This is a SEPARATE AMQP vhost from `broker:` above: member/worker inboxes always stay on the
|
||||
# per-fleet `broker:` vhost, and this vhost carries only leader-to-leader traffic, so two fleets
|
||||
# whose members must never see each other can still share one coordination vhost for their leads.
|
||||
# uriEnv → name of a host env var holding the coordination AMQP URI, same convention as
|
||||
# broker.uriEnv (keeps the credential out of fleetd.yaml). Wins over `uri` when set.
|
||||
# selfId → this daemon's own lead coord-id — the name its mailbox is owned under
|
||||
# (lead.<selfId>.inbox), e.g. "mac-opus" or "fleet01-lead". Must be globally unique
|
||||
# across every daemon sharing this vhost.
|
||||
# prefetch → consumer basicQos, capping how many unacked messages the mailbox holds in-heap.
|
||||
# Default 32 when omitted.
|
||||
# coordinator:
|
||||
# uriEnv: LEAD_COORD_URI
|
||||
# selfId: mac-opus
|
||||
# prefetch: 32
|
||||
|
||||
# Active push-to-primary (CB-307 Stage 3). When a worker reply lands with no open fleet_send,
|
||||
# the ReplyPushLoop injects a *drain nudge* (never the payload) into the primary's own herdr
|
||||
# pane — status-gated (only when injectable, never mid-turn) and bounded. Ack = drain: the loop
|
||||
|
||||
@@ -6,6 +6,7 @@ import com.fasterxml.jackson.core.JsonToken;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
|
||||
import dev.ltms.fleet.msg.AmqpReplyInbox;
|
||||
import dev.ltms.fleet.msg.LeadMailbox;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import org.slf4j.Logger;
|
||||
@@ -69,6 +70,11 @@ import java.util.Set;
|
||||
* @param memberCredentials deny-by-default policy (CB-596) for which of the operator's own host
|
||||
* credentials a spawned member's pane inherits. {@code null} (the block
|
||||
* omitted) blocks nothing — see {@link MemberCredentials}.
|
||||
* @param coordinator shared cross-host leader coordination broker: a SEPARATE AMQP vhost from
|
||||
* {@link #broker} used only for lead-to-lead traffic (member/worker inboxes
|
||||
* stay on {@code broker}'s vhost). {@code null} → no lead mailbox is opened.
|
||||
* Config parsing + accessors only — nothing here wires it into a live
|
||||
* {@code LeadMailbox}; that is a separate ticket. See {@link Coordinator}.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record FleetConfig(
|
||||
@@ -89,7 +95,20 @@ public record FleetConfig(
|
||||
Auth auth,
|
||||
ConfigReload configReload,
|
||||
Integer quarantineCooldownSeconds,
|
||||
MemberCredentials memberCredentials) {
|
||||
MemberCredentials memberCredentials,
|
||||
Coordinator coordinator) {
|
||||
|
||||
/** Back-compat form before the {@code coordinator:} block was added. */
|
||||
public FleetConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
|
||||
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
|
||||
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
|
||||
ConfigReload configReload, Integer quarantineCooldownSeconds,
|
||||
MemberCredentials memberCredentials) {
|
||||
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
|
||||
configReload, quarantineCooldownSeconds, memberCredentials, null);
|
||||
}
|
||||
|
||||
/** Back-compat form before the CB-596 {@code memberCredentials:} block was added. */
|
||||
public FleetConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
|
||||
@@ -650,6 +669,82 @@ public record FleetConfig(
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Shared cross-host leader coordination broker: a {@link LeadMailbox} lets two leads on
|
||||
* different daemons — possibly different hosts — exchange durable messages, which a herdr pane
|
||||
* injection (how {@code fleet_send} reaches a lead today) cannot do at all. Its mere presence is
|
||||
* config only in this ticket: nothing here opens a live {@code LeadMailbox} yet, that wiring is
|
||||
* a separate ticket.
|
||||
*
|
||||
* <p>Deliberately a SEPARATE vhost from {@link Broker}, not a reuse of it. {@link Broker} is
|
||||
* per-fleet — its queues are named by worker session id, and two fleets sharing one broker stay
|
||||
* isolated by vhost (see {@code Two fleets share one LavinMQ}). Leader coordination is meant to
|
||||
* cross exactly that boundary: two independently-owned fleets' leads talking to each other. Using
|
||||
* the same vhost would either leak member traffic across the fleet boundary this is meant to
|
||||
* cross, or force every fleet's members onto one shared vhost to get leader coordination — a
|
||||
* second, dedicated vhost keeps "member inboxes stay fleet-local" true while still letting leads
|
||||
* reach across fleets.
|
||||
*
|
||||
* @param uri AMQP connection URI for the coordination vhost, e.g.
|
||||
* {@code amqp://guest:guest@127.0.0.1:5672/coord}. Blank/{@code null} ⇒ the
|
||||
* coordinator block is treated as unconfigured. Ignored when {@code uriEnv} is set.
|
||||
* @param uriEnv name of a host env var holding the AMQP URI, same convention as
|
||||
* {@link Broker#uriEnv()} — keeps the credential out of the config file. Wins
|
||||
* over {@code uri} whenever set. Blank/{@code null} ⇒ ignored.
|
||||
* @param selfId this daemon's own lead coord-id — the name its {@code LeadMailbox} is owned
|
||||
* under ({@code lead.<selfId>.inbox}), e.g. {@code "mac-opus"}. Blank/
|
||||
* {@code null} ⇒ kept as {@code null} (no self id configured).
|
||||
* @param prefetch the consumer's {@code basicQos} prefetch count. {@code null}/non-positive ⇒
|
||||
* {@link LeadMailbox#DEFAULT_PREFETCH}.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record Coordinator(String uri, String uriEnv, String selfId, Integer prefetch) {
|
||||
|
||||
public Coordinator {
|
||||
selfId = (selfId == null || selfId.isBlank()) ? null : selfId;
|
||||
}
|
||||
|
||||
/** True when a {@code uriEnv} is configured by name, whether or not its variable resolves. */
|
||||
public boolean hasUriEnv() {
|
||||
return uriEnv != null && !uriEnv.isBlank();
|
||||
}
|
||||
|
||||
/**
|
||||
* True when a usable coordination broker URI is configured (an empty block does not enable
|
||||
* it). Honors {@code uriEnv} first: if it names a variable that is unset or blank, the
|
||||
* coordinator is <em>not</em> configured — a bare {@code uri} is only consulted when no
|
||||
* {@code uriEnv} is set.
|
||||
*/
|
||||
public boolean isConfigured() {
|
||||
return effectiveUri() != null;
|
||||
}
|
||||
|
||||
/**
|
||||
* The effective AMQP URI to connect with. {@code uriEnv} wins when set (both over
|
||||
* {@code uri} and alone). When {@code uriEnv} names a variable that is unset or blank,
|
||||
* returns {@code null} rather than falling back to {@code uri} — an operator who moved to
|
||||
* the secret store must not silently drop back onto a stale clear-text URI. Returns the
|
||||
* literal {@code uri} when no {@code uriEnv} is configured.
|
||||
*/
|
||||
public String effectiveUri() {
|
||||
return effectiveUri(System.getenv());
|
||||
}
|
||||
|
||||
/** As {@link #effectiveUri()}, reading the variable value from {@code env} (the injection seam). */
|
||||
public String effectiveUri(Map<String, String> env) {
|
||||
if (hasUriEnv()) {
|
||||
String value = env.get(uriEnv);
|
||||
return (value != null && !value.isBlank()) ? value : null;
|
||||
}
|
||||
return (uri != null && !uri.isBlank()) ? uri : null;
|
||||
}
|
||||
|
||||
/** The prefetch to use, defaulting to {@link LeadMailbox#DEFAULT_PREFETCH} when unset. */
|
||||
public int prefetchOrDefault() {
|
||||
return (prefetch != null && prefetch > 0) ? prefetch : LeadMailbox.DEFAULT_PREFETCH;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Optional pinned primary terminal config (CB-307). When present with a non-blank
|
||||
* {@code terminal}, the bridge uses this as the primary's herdr identity instead of
|
||||
@@ -1184,7 +1279,7 @@ public record FleetConfig(
|
||||
"bind", "herdrSocket", "profiles", "guard", "worktreeRoot",
|
||||
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
|
||||
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
|
||||
"memberCredentials");
|
||||
"memberCredentials", "coordinator");
|
||||
|
||||
/** Load and validate config from {@code path}. */
|
||||
public static FleetConfig load(Path path) {
|
||||
@@ -1753,9 +1848,11 @@ public record FleetConfig(
|
||||
// and an upgrade must not change what a running deployment's members inherit.
|
||||
MemberCredentials mc = memberCredentials != null ? memberCredentials
|
||||
: new MemberCredentials(null, List.of(), List.of());
|
||||
// coordinator is left as-is, like broker/primary above: null keeps no LeadMailbox opened,
|
||||
// and this ticket's Coordinator is config-only anyway (nothing yet reads it at startup).
|
||||
return new FleetConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
|
||||
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
|
||||
quarantineCooldown, mc);
|
||||
quarantineCooldown, mc, coordinator);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -0,0 +1,422 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.rabbitmq.client.AMQP;
|
||||
import com.rabbitmq.client.Channel;
|
||||
import com.rabbitmq.client.Connection;
|
||||
import com.rabbitmq.client.ConnectionFactory;
|
||||
import com.rabbitmq.client.DeliverCallback;
|
||||
import com.rabbitmq.client.Recoverable;
|
||||
import com.rabbitmq.client.RecoveryListener;
|
||||
import com.rabbitmq.client.Return;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.NavigableMap;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ConcurrentSkipListMap;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
|
||||
/**
|
||||
* Durable, AMQP-backed mailbox for lead-to-lead messages across daemons — including daemons on
|
||||
* different hosts, where a herdr pane injection (how {@code fleet_send} reaches a lead today)
|
||||
* cannot reach at all. The broker is the only medium two independently-owned daemons share, which
|
||||
* is exactly why {@link AmqpReplyInbox}'s javadoc already calls out "one gateway may publish to an
|
||||
* agent owned by another gateway" (CB-308 federation) as the reason {@code publish} and
|
||||
* {@code own}/consume are separate operations there — this class leans on the same split.
|
||||
*
|
||||
* <p><strong>Single-target, unlike {@link AmqpReplyInbox}.</strong> {@code AmqpReplyInbox}
|
||||
* multiplexes many workers' reply queues under one gateway connection. A {@code LeadMailbox}
|
||||
* instance is simpler: it owns exactly <em>one</em> queue — this daemon's own
|
||||
* {@code lead.<selfCoordId>.inbox} — declared and consumed the moment it is constructed. There is
|
||||
* no {@code own}/{@code release} pair to call separately; a daemon either runs a {@code LeadMailbox}
|
||||
* for its own coord-id, or it does not run one at all.
|
||||
*
|
||||
* <p><strong>Consume-and-hold with deferred manual ack</strong> — same mapping as
|
||||
* {@code AmqpReplyInbox}. The constructor declares the durable queue and starts a manual-ack
|
||||
* consumer that pulls persistent messages into an in-memory {@code held} map (keyed by
|
||||
* {@link LeadMessage#msgId()}) but does not ack them. {@link #peek} returns a non-destructive
|
||||
* snapshot; {@link #ack} acks the broker delivery-tag and drops the entry. A message that is never
|
||||
* acked (a crash, a bounce) survives — the broker redelivers it to the next connection that owns
|
||||
* the queue.
|
||||
*
|
||||
* <p><strong>Publishing does not imply owning.</strong> {@link #publish} sends to
|
||||
* {@code lead.<toCoordId>.inbox} over a dedicated confirm-mode channel; it never declares that
|
||||
* queue as owned and never attaches a consumer to it. A sender that has never opened its own
|
||||
* {@code LeadMailbox} for {@code toCoordId} can still publish to it, exactly as CB-308 federation
|
||||
* requires. Publish blocks for the broker's publisher confirm (persistent delivery, {@code
|
||||
* mandatory=true}) and throws {@link IllegalStateException} on an unroutable return, a nack, or a
|
||||
* timeout — the caller must not report success for a black-holed message.
|
||||
*
|
||||
* <p><strong>Recovery.</strong> The connection is opened with automatic + topology recovery
|
||||
* enabled, mirroring {@code AmqpReplyInbox}: on reconnect the broker hands out fresh delivery tags,
|
||||
* so the held snapshot is cleared (dedup by {@code msgId} still prevents any double-queue on
|
||||
* redelivery) and any publish still awaiting its confirm is failed rather than left to idle out
|
||||
* the confirm timeout against a sequence number that means nothing on the new channel.
|
||||
*/
|
||||
public final class LeadMailbox implements AutoCloseable {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(LeadMailbox.class);
|
||||
|
||||
private static final String QUEUE_PREFIX = "lead.";
|
||||
private static final String QUEUE_SUFFIX = ".inbox";
|
||||
|
||||
/** The prefetch used when a caller does not pass an explicit value to {@link #open(String, String, int)}. */
|
||||
public static final int DEFAULT_PREFETCH = 32;
|
||||
|
||||
/** How long {@link #publish} waits for its publisher confirm before failing the call. */
|
||||
private static final long CONFIRM_TIMEOUT_MS = 10_000L;
|
||||
|
||||
private static final ObjectMapper MAPPER = new ObjectMapper();
|
||||
|
||||
private final Connection connection;
|
||||
private final String selfCoordId;
|
||||
|
||||
private final Channel channel;
|
||||
/** All consume-channel operations (declare/consume/ack) serialize on this — a Channel is not thread-safe. */
|
||||
private final Object channelLock = new Object();
|
||||
/** msgId → held delivery, for this mailbox's own queue only (there is exactly one). */
|
||||
private final LinkedHashMap<String, Held> held = new LinkedHashMap<>();
|
||||
|
||||
/**
|
||||
* A dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume + ack)
|
||||
* so a publish confirm round trip never blocks under {@link #channelLock} and stalls an ack.
|
||||
*/
|
||||
private final Channel publishChannel;
|
||||
private final Object publishChannelLock = new Object();
|
||||
/** In-flight publishes awaiting their confirm, keyed by the publish channel's sequence number. */
|
||||
private final ConcurrentSkipListMap<Long, Pending> pendingBySeq = new ConcurrentSkipListMap<>();
|
||||
/** The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery tag. */
|
||||
private final ConcurrentHashMap<String, Pending> pendingByMsgId = new ConcurrentHashMap<>();
|
||||
|
||||
/** A message pulled off the broker but not yet acked: its delivery-tag plus the deserialized envelope. */
|
||||
private record Held(long deliveryTag, LeadMessage message) {}
|
||||
|
||||
/** A publish awaiting its confirm; {@link #returned} records whether the broker already returned it. */
|
||||
private static final class Pending {
|
||||
final String msgId;
|
||||
final CompletableFuture<Void> confirmed = new CompletableFuture<>();
|
||||
volatile boolean returned;
|
||||
|
||||
Pending(String msgId) {
|
||||
this.msgId = msgId;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Connect to {@code uri} (the shared cross-host coordination vhost, e.g.
|
||||
* {@code amqp://guest:guest@127.0.0.1:5672/coord}) and own {@code selfCoordId}'s mailbox, with
|
||||
* {@link #DEFAULT_PREFETCH}.
|
||||
*/
|
||||
public static LeadMailbox open(String uri, String selfCoordId) {
|
||||
return open(uri, selfCoordId, DEFAULT_PREFETCH);
|
||||
}
|
||||
|
||||
/** As {@link #open(String, String)}, with an explicit consumer prefetch. */
|
||||
public static LeadMailbox open(String uri, String selfCoordId, int prefetch) {
|
||||
try {
|
||||
ConnectionFactory factory = new ConnectionFactory();
|
||||
factory.setUri(uri);
|
||||
// Self-heal transient blips; topology recovery re-declares the queue and re-attaches the consumer.
|
||||
factory.setAutomaticRecoveryEnabled(true);
|
||||
factory.setTopologyRecoveryEnabled(true);
|
||||
return new LeadMailbox(factory.newConnection("bridged-lead-mailbox"), selfCoordId, prefetch);
|
||||
} catch (Exception e) {
|
||||
throw new IllegalStateException("cannot connect to AMQP coordination broker at " + uri, e);
|
||||
}
|
||||
}
|
||||
|
||||
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for tests). */
|
||||
LeadMailbox(Connection connection, String selfCoordId) {
|
||||
this(connection, selfCoordId, DEFAULT_PREFETCH);
|
||||
}
|
||||
|
||||
/** As above, with an explicit prefetch (injection seam for tests). */
|
||||
LeadMailbox(Connection connection, String selfCoordId, int prefetch) {
|
||||
this.connection = connection;
|
||||
this.selfCoordId = selfCoordId;
|
||||
try {
|
||||
this.channel = connection.createChannel();
|
||||
// Bound the held backlog — must be set before basicConsume.
|
||||
this.channel.basicQos(prefetch);
|
||||
this.publishChannel = connection.createChannel();
|
||||
this.publishChannel.confirmSelect();
|
||||
this.publishChannel.addReturnListener(this::onReturn);
|
||||
this.publishChannel.addConfirmListener(this::onAck, this::onNack);
|
||||
own();
|
||||
} catch (IOException e) {
|
||||
throw new IllegalStateException("cannot open AMQP channel", e);
|
||||
}
|
||||
// On automatic recovery the broker redelivers unacked messages with FRESH delivery-tags; the
|
||||
// tags we were holding are now stale. Drop the held snapshot so the re-attached consumer
|
||||
// repopulates it with valid tags (dedup by msgId still prevents any double-queue). Any publish
|
||||
// confirm still in flight when the connection dropped is equally stale — fail it now rather
|
||||
// than let it silently ride out CONFIRM_TIMEOUT_MS.
|
||||
if (connection instanceof Recoverable recoverable) {
|
||||
recoverable.addRecoveryListener(new RecoveryListener() {
|
||||
@Override
|
||||
public void handleRecovery(Recoverable recoverable) {
|
||||
synchronized (held) {
|
||||
held.clear();
|
||||
}
|
||||
failPendingPublishesOnRecovery();
|
||||
log.info("AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery");
|
||||
}
|
||||
|
||||
@Override
|
||||
public void handleRecoveryStarted(Recoverable recoverable) {
|
||||
// no-op: we act once recovery completes
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
/** Declare + consume this daemon's own {@code lead.<selfCoordId>.inbox}. Called once, at construction. */
|
||||
private void own() throws IOException {
|
||||
String queue = queueName(selfCoordId);
|
||||
synchronized (channelLock) {
|
||||
channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
|
||||
channel.basicConsume(queue, false, deliverCallback(), _ -> { }); // autoAck=false: manual ack
|
||||
}
|
||||
log.debug("lead mailbox owns queue {} for coord-id {}", queue, selfCoordId);
|
||||
}
|
||||
|
||||
/**
|
||||
* Publish {@code msg} to {@code toCoordId}'s mailbox and block until the broker's publisher
|
||||
* confirm for it lands. Does <em>not</em> imply owning or consuming {@code toCoordId}'s queue.
|
||||
* Throws {@link IllegalStateException} if the message is returned as unroutable, nacked, or not
|
||||
* confirmed within {@link #CONFIRM_TIMEOUT_MS} — the caller must treat that as a failed publish,
|
||||
* not a lost-and-forgotten one.
|
||||
*/
|
||||
public void publish(String toCoordId, LeadMessage msg) {
|
||||
byte[] body;
|
||||
try {
|
||||
body = MAPPER.writeValueAsBytes(msg);
|
||||
} catch (JsonProcessingException e) {
|
||||
throw new IllegalStateException("cannot serialize lead message " + msg.msgId(), e);
|
||||
}
|
||||
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
|
||||
.messageId(msg.msgId())
|
||||
.deliveryMode(2) // persistent — survives a broker restart
|
||||
.contentType("application/json")
|
||||
.build();
|
||||
Pending pending = new Pending(msg.msgId());
|
||||
long seq;
|
||||
synchronized (publishChannelLock) {
|
||||
seq = publishChannel.getNextPublishSeqNo();
|
||||
pendingBySeq.put(seq, pending);
|
||||
pendingByMsgId.put(msg.msgId(), pending);
|
||||
try {
|
||||
publishChannel.basicPublish("", queueName(toCoordId), true, props, body);
|
||||
} catch (IOException e) {
|
||||
pendingBySeq.remove(seq, pending);
|
||||
pendingByMsgId.remove(msg.msgId(), pending);
|
||||
throw new IllegalStateException("cannot publish lead message to " + queueName(toCoordId), e);
|
||||
}
|
||||
}
|
||||
try {
|
||||
pending.confirmed.get(CONFIRM_TIMEOUT_MS, TimeUnit.MILLISECONDS);
|
||||
} catch (ExecutionException e) {
|
||||
Throwable cause = e.getCause();
|
||||
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
|
||||
} catch (TimeoutException e) {
|
||||
throw new IllegalStateException("publish confirm for lead message " + msg.msgId() + " to "
|
||||
+ queueName(toCoordId) + " timed out after " + CONFIRM_TIMEOUT_MS
|
||||
+ "ms — broker may be unreachable or overloaded", e);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new IllegalStateException("interrupted awaiting publish confirm for " + msg.msgId(), e);
|
||||
} finally {
|
||||
pendingBySeq.remove(seq, pending);
|
||||
pendingByMsgId.remove(msg.msgId(), pending);
|
||||
}
|
||||
}
|
||||
|
||||
/** Non-destructive FIFO snapshot of this mailbox's currently-held messages. */
|
||||
public List<LeadMessage> peek() {
|
||||
synchronized (held) {
|
||||
return held.values().stream().map(Held::message).toList();
|
||||
}
|
||||
}
|
||||
|
||||
/** Convenience: {@link #peek} the current snapshot, then {@link #ack} every message in it. */
|
||||
public List<LeadMessage> drain() {
|
||||
List<LeadMessage> snapshot = peek();
|
||||
snapshot.forEach(m -> ack(m.msgId()));
|
||||
return snapshot;
|
||||
}
|
||||
|
||||
/** Remove the held message {@code msgId} and ack it on the broker. No-op if not held. */
|
||||
public void ack(String msgId) {
|
||||
Held h;
|
||||
synchronized (held) {
|
||||
h = held.remove(msgId);
|
||||
}
|
||||
if (h == null) {
|
||||
return; // never held (or already acked) — no-op
|
||||
}
|
||||
try {
|
||||
synchronized (channelLock) {
|
||||
channel.basicAck(h.deliveryTag(), false);
|
||||
}
|
||||
} catch (IOException e) {
|
||||
// Ack didn't reach the broker: restore the entry so a later ack (or a redelivery after
|
||||
// reconnect) can retry. Keeps the at-least-once contract — a message is never silently lost.
|
||||
synchronized (held) {
|
||||
held.putIfAbsent(msgId, h);
|
||||
}
|
||||
throw new IllegalStateException("cannot ack lead message " + msgId, e);
|
||||
}
|
||||
}
|
||||
|
||||
private DeliverCallback deliverCallback() {
|
||||
return (_, delivery) -> {
|
||||
long tag = delivery.getEnvelope().getDeliveryTag();
|
||||
LeadMessage msg;
|
||||
try {
|
||||
msg = MAPPER.readValue(delivery.getBody(), LeadMessage.class);
|
||||
} catch (IOException e) {
|
||||
// A malformed body can never be dedup-keyed or handed to a caller; ack it so the
|
||||
// broker does not redeliver it forever, and log loudly since this should never happen
|
||||
// for a producer that only ever calls publish(String, LeadMessage).
|
||||
log.warn("dropping malformed lead-mailbox delivery (tag {}): {}", tag, e.toString());
|
||||
synchronized (channelLock) {
|
||||
channel.basicAck(tag, false);
|
||||
}
|
||||
return;
|
||||
}
|
||||
boolean duplicate;
|
||||
synchronized (held) {
|
||||
if (held.containsKey(msg.msgId())) {
|
||||
duplicate = true;
|
||||
} else {
|
||||
held.put(msg.msgId(), new Held(tag, msg));
|
||||
duplicate = false;
|
||||
}
|
||||
}
|
||||
if (duplicate) {
|
||||
// Redelivered duplicate: ack the new tag and drop it so the broker stops resending.
|
||||
synchronized (channelLock) {
|
||||
channel.basicAck(tag, false);
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
/** Broker return for an unroutable {@code mandatory} publish — arrives BEFORE its confirm. */
|
||||
private void onReturn(Return r) {
|
||||
String msgId = r.getProperties() == null ? null : r.getProperties().getMessageId();
|
||||
Pending pending = msgId == null ? null : pendingByMsgId.get(msgId);
|
||||
if (pending != null) {
|
||||
pending.returned = true;
|
||||
} else {
|
||||
log.warn("AMQP return for lead message {} (routingKey={}, {} {}) with no matching in-flight publish"
|
||||
+ " — already resolved by a prior confirm", msgId, r.getRoutingKey(), r.getReplyCode(),
|
||||
r.getReplyText());
|
||||
}
|
||||
}
|
||||
|
||||
private void onAck(long seq, boolean multiple) {
|
||||
resolveConfirm(seq, multiple, true);
|
||||
}
|
||||
|
||||
private void onNack(long seq, boolean multiple) {
|
||||
resolveConfirm(seq, multiple, false);
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve every pending publish covered by this confirm (a single seq, or — {@code multiple} —
|
||||
* every seq up to and including it). Checks {@link Pending#returned} at confirm time: since the
|
||||
* broker's return for an unroutable message always precedes its confirm, an ack that arrives after
|
||||
* a return means "confirmed but never routed", not "durably queued".
|
||||
*/
|
||||
private void resolveConfirm(long seq, boolean multiple, boolean ack) {
|
||||
NavigableMap<Long, Pending> covered = multiple
|
||||
? pendingBySeq.headMap(seq, true)
|
||||
: pendingBySeq.subMap(seq, true, seq, true);
|
||||
for (var it = covered.entrySet().iterator(); it.hasNext(); ) {
|
||||
Pending pending = it.next().getValue();
|
||||
it.remove();
|
||||
pendingByMsgId.remove(pending.msgId, pending);
|
||||
if (ack && !pending.returned) {
|
||||
pending.confirmed.complete(null);
|
||||
} else if (ack) {
|
||||
pending.confirmed.completeExceptionally(new IllegalStateException(
|
||||
"lead message " + pending.msgId + " was returned as unroutable (mailbox not owned)"));
|
||||
} else {
|
||||
pending.confirmed.completeExceptionally(new IllegalStateException(
|
||||
"broker nacked publish of lead message " + pending.msgId));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Fail every publish still awaiting its confirm — their sequence numbers are stale after
|
||||
* recovery. Guarded by {@link #publishChannelLock}, the same lock {@link #publish} holds while
|
||||
* it takes its sequence number and registers its {@link Pending} — see
|
||||
* {@code AmqpReplyInbox.failPendingPublishesOnRecovery}'s javadoc for the full race analysis this
|
||||
* mirrors. Package-private only so a unit test can drive it directly without a live broker
|
||||
* reconnect.
|
||||
*/
|
||||
void failPendingPublishesOnRecovery() {
|
||||
synchronized (publishChannelLock) {
|
||||
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
|
||||
Pending pending = it.next().getValue();
|
||||
it.remove();
|
||||
pendingByMsgId.remove(pending.msgId, pending);
|
||||
pending.confirmed.completeExceptionally(new IllegalStateException(
|
||||
"AMQP connection recovered mid-publish; confirm status of lead message "
|
||||
+ pending.msgId + " is unknown"));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Fail every publish still awaiting its confirm with a clear, immediate error instead of leaving
|
||||
* it to time out after {@link #CONFIRM_TIMEOUT_MS} once the channels are closed underneath it.
|
||||
*/
|
||||
private void failPendingPublishesOnClose() {
|
||||
synchronized (publishChannelLock) {
|
||||
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
|
||||
Pending pending = it.next().getValue();
|
||||
it.remove();
|
||||
pendingByMsgId.remove(pending.msgId, pending);
|
||||
pending.confirmed.completeExceptionally(new IllegalStateException(
|
||||
"lead mailbox closed while publish of lead message " + pending.msgId
|
||||
+ " was still awaiting its confirm"));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** The durable queue name a coord-id's mailbox lives on: {@code lead.<coordId>.inbox}. */
|
||||
public static String queueName(String coordId) {
|
||||
return QUEUE_PREFIX + coordId + QUEUE_SUFFIX;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
failPendingPublishesOnClose();
|
||||
try {
|
||||
channel.close();
|
||||
} catch (Exception e) {
|
||||
log.debug("AMQP lead mailbox channel close: {}", e.toString());
|
||||
}
|
||||
try {
|
||||
publishChannel.close();
|
||||
} catch (Exception e) {
|
||||
log.debug("AMQP lead mailbox publish channel close: {}", e.toString());
|
||||
}
|
||||
try {
|
||||
connection.close();
|
||||
} catch (Exception e) {
|
||||
log.debug("AMQP lead mailbox connection close: {}", e.toString());
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,22 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
/**
|
||||
* Wire envelope for a lead-to-lead message carried over {@link LeadMailbox}.
|
||||
*
|
||||
* <p>Unlike {@link ReplyInbox.InboxMessage} (a worker→primary reply, addressed only by the single
|
||||
* gateway that owns the worker), a lead message crosses independently-owned daemons — possibly on
|
||||
* different hosts — so it carries an explicit sender ({@code from}) as well as the recipient
|
||||
* ({@code to}): the recipient needs the sender's coord-id to reply back.
|
||||
*
|
||||
* <p>{@code from} and {@code to} are globally-unique lead coordination ids (e.g. {@code "mac-opus"},
|
||||
* {@code "fleet01-lead"}) — NOT herdr terminal ids. A herdr terminal id is meaningful only on the
|
||||
* host that owns it, so it cannot address a lead running on another daemon; a coord-id is chosen
|
||||
* by configuration ({@code coordinator.selfId}) precisely so it means the same thing everywhere.
|
||||
*
|
||||
* @param msgId idempotency id; a redelivered duplicate (at-least-once delivery) is deduped on this
|
||||
* @param from the sending lead's coord-id
|
||||
* @param to the receiving lead's coord-id — identifies the mailbox this message is held on
|
||||
* @param content the message text
|
||||
*/
|
||||
public record LeadMessage(String msgId, String from, String to, String content) {
|
||||
}
|
||||
@@ -1,6 +1,7 @@
|
||||
package dev.ltms.fleet.config;
|
||||
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.msg.LeadMailbox;
|
||||
import dev.ltms.fleet.peer.MemberRole;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
@@ -932,6 +933,66 @@ class FleetConfigTest {
|
||||
assertFalse(cfg.broker().isConfigured(), "an empty uri must not enable AMQP");
|
||||
}
|
||||
|
||||
@Test
|
||||
void absentCoordinatorBlockLeavesCoordinatorNull(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-coordinator.yaml");
|
||||
Files.writeString(f, "bind:\n port: 8080\n");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertNull(cfg.coordinator(), "no coordinator: block → null → no lead mailbox is opened");
|
||||
}
|
||||
|
||||
@Test
|
||||
void coordinatorBlockParses(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("coordinator.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
coordinator:
|
||||
uri: amqp://guest:guest@127.0.0.1:5672/coord
|
||||
selfId: mac-opus
|
||||
prefetch: 16
|
||||
""");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertNotNull(cfg.coordinator());
|
||||
assertTrue(cfg.coordinator().isConfigured(), "a non-blank uri enables the coordinator");
|
||||
assertEquals("amqp://guest:guest@127.0.0.1:5672/coord", cfg.coordinator().uri());
|
||||
assertEquals("mac-opus", cfg.coordinator().selfId());
|
||||
assertEquals(16, cfg.coordinator().prefetchOrDefault());
|
||||
}
|
||||
|
||||
@Test
|
||||
void coordinatorBlockWithBlankUriStaysUnconfigured(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("coordinator-blank.yaml");
|
||||
Files.writeString(f, "bind:\n port: 8080\ncoordinator:\n uri: \"\"\n");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertNotNull(cfg.coordinator());
|
||||
assertFalse(cfg.coordinator().isConfigured(), "an empty uri must not enable the coordinator");
|
||||
assertNull(cfg.coordinator().selfId(), "a blank/absent selfId stays null, never coerced to empty");
|
||||
}
|
||||
|
||||
@Test
|
||||
void coordinatorEffectiveUriHonorsUriEnv() {
|
||||
FleetConfig.Coordinator withEnv = new FleetConfig.Coordinator(
|
||||
"amqp://stale-clear-text@127.0.0.1:5672/coord", "LEAD_COORD_URI", "fleet01-lead", null);
|
||||
|
||||
assertEquals("amqp://from-env@127.0.0.1:5672/coord",
|
||||
withEnv.effectiveUri(Map.of("LEAD_COORD_URI", "amqp://from-env@127.0.0.1:5672/coord")),
|
||||
"uriEnv wins over a literal uri when its variable resolves");
|
||||
assertNull(withEnv.effectiveUri(Map.of()),
|
||||
"an unset uriEnv variable must not fall back to the literal uri");
|
||||
assertNull(withEnv.effectiveUri(Map.of("LEAD_COORD_URI", " ")),
|
||||
"a blank uriEnv variable must not fall back to the literal uri");
|
||||
|
||||
FleetConfig.Coordinator noEnv = new FleetConfig.Coordinator(
|
||||
"amqp://guest:guest@127.0.0.1:5672/coord", null, null, null);
|
||||
assertEquals("amqp://guest:guest@127.0.0.1:5672/coord", noEnv.effectiveUri(Map.of()),
|
||||
"the literal uri is used when no uriEnv is configured");
|
||||
assertEquals(LeadMailbox.DEFAULT_PREFETCH, noEnv.prefetchOrDefault());
|
||||
}
|
||||
|
||||
@Test
|
||||
void absentPrimaryBlockLeavesPrimaryNull(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-primary.yaml");
|
||||
|
||||
@@ -0,0 +1,173 @@
|
||||
package dev.ltms.fleet.msg;
|
||||
|
||||
import org.junit.jupiter.api.BeforeAll;
|
||||
import org.junit.jupiter.api.Tag;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.testcontainers.containers.RabbitMQContainer;
|
||||
import org.testcontainers.junit.jupiter.Testcontainers;
|
||||
import org.testcontainers.utility.DockerImageName;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* Contract test for {@link LeadMailbox} against a REAL broker — same approach as
|
||||
* {@code AmqpReplyInboxContractTest}, which this mirrors: a Testcontainers RabbitMQ locally, or an
|
||||
* externally-provisioned broker in CI via {@code AMQP_URI}. Tagged {@code contract} so it is
|
||||
* excluded from {@code mvn test}/{@code mvn clean install} (which stay hermetic and need no
|
||||
* Docker); run it with Docker present via {@code mvn test -Pcontract}.
|
||||
*
|
||||
* <p>Proves the mechanism this ticket adds: a {@link LeadMessage} published to
|
||||
* {@code lead.<to>.inbox} is received with {@code from}/{@code to}/{@code content} intact, and
|
||||
* {@link LeadMailbox#ack} removes it — the same publish→peek→ack roundtrip
|
||||
* {@code AmqpReplyInboxContractTest} proves for {@link AmqpReplyInbox}, adapted to this class's
|
||||
* single-owned-mailbox shape (no {@code own}/{@code release} — the mailbox for {@code selfCoordId}
|
||||
* is owned the moment {@link LeadMailbox#open} returns).
|
||||
*/
|
||||
@Tag("contract")
|
||||
// disabledWithoutDocker=false: on the CI path (AMQP_URI set) no container is started and the class
|
||||
// must still run against the external broker even though the runner has no Docker.
|
||||
@Testcontainers(disabledWithoutDocker = false)
|
||||
class LeadMailboxTest {
|
||||
|
||||
private static final String EXTERNAL_URI = System.getenv("AMQP_URI");
|
||||
|
||||
private static final RabbitMQContainer BROKER =
|
||||
new RabbitMQContainer(DockerImageName.parse("rabbitmq:3.13-management"));
|
||||
|
||||
private static final AtomicLong SEQ = new AtomicLong();
|
||||
|
||||
// No @Container: the JUnit 5 extension would force-start it even when AMQP_URI is set. Start it
|
||||
// manually only on the local (no-external-broker) path; Ryuk reaps it on JVM exit.
|
||||
@BeforeAll
|
||||
static void startBrokerUnlessExternal() {
|
||||
if (EXTERNAL_URI == null) {
|
||||
BROKER.start();
|
||||
}
|
||||
}
|
||||
|
||||
private static String uri() {
|
||||
if (EXTERNAL_URI != null) {
|
||||
return EXTERNAL_URI;
|
||||
}
|
||||
// No trailing slash: an empty path is vhost "", which does not exist — omitting it selects
|
||||
// the default vhost "/".
|
||||
return "amqp://guest:guest@" + BROKER.getHost() + ":" + BROKER.getAmqpPort();
|
||||
}
|
||||
|
||||
/** A fresh coord-id per test run so parallel/repeat runs never collide on the same queue. */
|
||||
private static String coordId(String prefix) {
|
||||
return prefix + "-" + System.nanoTime() + "-" + SEQ.incrementAndGet();
|
||||
}
|
||||
|
||||
@Test
|
||||
void publishThenPeekThenAckRoundTrip() throws Exception {
|
||||
String to = coordId("lead-to");
|
||||
String from = "lead-from";
|
||||
try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) {
|
||||
LeadMessage sent = new LeadMessage("m1", from, to, "hello peer lead");
|
||||
inbox.publish(to, sent);
|
||||
|
||||
List<LeadMessage> got = awaitPeek(inbox);
|
||||
assertEquals(1, got.size(), "the published message should be held for drain");
|
||||
assertEquals("m1", got.getFirst().msgId());
|
||||
assertEquals(from, got.getFirst().from());
|
||||
assertEquals(to, got.getFirst().to());
|
||||
assertEquals("hello peer lead", got.getFirst().content());
|
||||
|
||||
inbox.ack("m1");
|
||||
assertTrue(inbox.peek().isEmpty(), "an acked message is dropped");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void duplicateMsgIdIsNotDoubleQueued() throws Exception {
|
||||
String to = coordId("lead-dedup");
|
||||
try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) {
|
||||
inbox.publish(to, new LeadMessage("dup", "lead-from", to, "first"));
|
||||
awaitPeek(inbox);
|
||||
inbox.publish(to, new LeadMessage("dup", "lead-from", to, "second")); // same msgId — no-op
|
||||
|
||||
Thread.sleep(500); // give any erroneous second delivery time to land
|
||||
List<LeadMessage> got = inbox.peek();
|
||||
assertEquals(1, got.size(), "a repeated msgId must not double-queue");
|
||||
assertEquals("first", got.getFirst().content(), "the first payload wins");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void unackedMessageSurvivesRestartAndIsRedelivered() throws Exception {
|
||||
String to = coordId("lead-durable");
|
||||
|
||||
// First "process life": publish, see it held, but crash before acking.
|
||||
try (LeadMailbox first = LeadMailbox.open(uri(), to)) {
|
||||
first.publish(to, new LeadMessage("persist-1", "lead-from", to, "survive me"));
|
||||
assertEquals(1, awaitPeek(first).size());
|
||||
// no ack — simulate a java -jar bounce with the message still pending
|
||||
}
|
||||
|
||||
// Second "process life": a fresh connection owning the same mailbox must be redelivered it.
|
||||
try (LeadMailbox second = LeadMailbox.open(uri(), to)) {
|
||||
List<LeadMessage> got = awaitPeek(second);
|
||||
assertEquals(1, got.size(), "an unacked persistent message is redelivered after restart");
|
||||
assertEquals("persist-1", got.getFirst().msgId());
|
||||
assertEquals("survive me", got.getFirst().content());
|
||||
|
||||
second.ack("persist-1");
|
||||
}
|
||||
|
||||
// Third life: once acked, it is gone for good.
|
||||
try (LeadMailbox third = LeadMailbox.open(uri(), to)) {
|
||||
Thread.sleep(500);
|
||||
assertTrue(third.peek().isEmpty(), "an acked message does not come back on the next restart");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void publishDoesNotRequireTheSenderToOwnTheTargetMailbox() throws Exception {
|
||||
// CB-308 federation: a sender that never opened its own LeadMailbox for `to` can still
|
||||
// publish to it — publish must not imply ownership. Only the owner ever consumes here.
|
||||
String to = coordId("lead-federated");
|
||||
String senderId = coordId("lead-sender");
|
||||
try (LeadMailbox owner = LeadMailbox.open(uri(), to);
|
||||
LeadMailbox sender = LeadMailbox.open(uri(), senderId)) {
|
||||
sender.publish(to, new LeadMessage("m1", senderId, to, "from a federated peer"));
|
||||
|
||||
List<LeadMessage> got = awaitPeek(owner);
|
||||
assertEquals(1, got.size(), "only the owner's mailbox should receive the message");
|
||||
assertEquals(senderId, got.getFirst().from());
|
||||
assertTrue(sender.peek().isEmpty(), "the sender must not also hold a copy — it never owns `to`");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void unroutablePublishReportsFailureNotSilentSuccess() throws Exception {
|
||||
// Publish to a coord-id whose mailbox was never opened by anyone: the queue is never
|
||||
// declared, so the default-exchange route to lead.<to>.inbox does not exist and the broker
|
||||
// must return the publish.
|
||||
String to = coordId("lead-nobody-home");
|
||||
try (LeadMailbox sender = LeadMailbox.open(uri(), coordId("lead-sender"))) {
|
||||
IllegalStateException ex = assertThrows(IllegalStateException.class,
|
||||
() -> sender.publish(to, new LeadMessage("m1", "lead-from", to, "nobody home")));
|
||||
assertTrue(ex.getMessage() != null && ex.getMessage().toLowerCase().contains("unroutable"),
|
||||
"expected an unroutable-publish failure, got: " + ex.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
/** Poll peek until at least one message is held, or ~10s elapse (broker delivery is async). */
|
||||
@SuppressWarnings("BusyWait")
|
||||
private static List<LeadMessage> awaitPeek(LeadMailbox inbox) throws InterruptedException {
|
||||
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
|
||||
List<LeadMessage> msgs = inbox.peek();
|
||||
while (msgs.isEmpty() && System.nanoTime() < deadline) {
|
||||
Thread.sleep(50);
|
||||
msgs = inbox.peek();
|
||||
}
|
||||
return msgs;
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user