From 7b98cca967cbc64e8fe27297b8a348c75093f594 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Mon, 24 Aug 2026 17:23:42 +0200 Subject: [PATCH] LeadMailbox: durable leader-to-leader mailbox over a shared coordination vhost MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Adds the broker-side mechanism for lead-to-lead messages across daemons/hosts (unit 1 of 2): a LeadMessage envelope carrying from/to coord-ids, an AMQP-backed LeadMailbox modeled closely on AmqpReplyInbox (consume-and-hold, deferred manual ack, confirm-mode publish, recovery handling), and a new optional coordinator: config block (separate vhost from broker:, leader traffic only). Config parsing + accessors only — FleetMcp/Fleetd/Injector/ MessageService and the send path are untouched; wiring is a separate ticket. --- bridged/fleetd.example.yaml | 17 + .../dev/ltms/fleet/config/FleetConfig.java | 103 ++++- .../java/dev/ltms/fleet/msg/LeadMailbox.java | 422 ++++++++++++++++++ .../java/dev/ltms/fleet/msg/LeadMessage.java | 22 + .../ltms/fleet/config/FleetConfigTest.java | 61 +++ .../dev/ltms/fleet/msg/LeadMailboxTest.java | 173 +++++++ 6 files changed, 795 insertions(+), 3 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java create mode 100644 bridged/src/main/java/dev/ltms/fleet/msg/LeadMessage.java create mode 100644 bridged/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java diff --git a/bridged/fleetd.example.yaml b/bridged/fleetd.example.yaml index ebfe523..ba58eb0 100644 --- a/bridged/fleetd.example.yaml +++ b/bridged/fleetd.example.yaml @@ -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..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 diff --git a/bridged/src/main/java/dev/ltms/fleet/config/FleetConfig.java b/bridged/src/main/java/dev/ltms/fleet/config/FleetConfig.java index 994d98b..c616d66 100644 --- a/bridged/src/main/java/dev/ltms/fleet/config/FleetConfig.java +++ b/bridged/src/main/java/dev/ltms/fleet/config/FleetConfig.java @@ -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 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 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. + * + *

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..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 not 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 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); } /** diff --git a/bridged/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java b/bridged/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java new file mode 100644 index 0000000..4ea4d6a --- /dev/null +++ b/bridged/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java @@ -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. + * + *

Single-target, unlike {@link AmqpReplyInbox}. {@code AmqpReplyInbox} + * multiplexes many workers' reply queues under one gateway connection. A {@code LeadMailbox} + * instance is simpler: it owns exactly one queue — this daemon's own + * {@code lead..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. + * + *

Consume-and-hold with deferred manual ack — 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. + * + *

Publishing does not imply owning. {@link #publish} sends to + * {@code lead..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. + * + *

Recovery. 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 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 pendingBySeq = new ConcurrentSkipListMap<>(); + /** The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery tag. */ + private final ConcurrentHashMap 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 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..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 not 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 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 drain() { + List 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 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..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()); + } + } +} diff --git a/bridged/src/main/java/dev/ltms/fleet/msg/LeadMessage.java b/bridged/src/main/java/dev/ltms/fleet/msg/LeadMessage.java new file mode 100644 index 0000000..4b7bc3d --- /dev/null +++ b/bridged/src/main/java/dev/ltms/fleet/msg/LeadMessage.java @@ -0,0 +1,22 @@ +package dev.ltms.fleet.msg; + +/** + * Wire envelope for a lead-to-lead message carried over {@link LeadMailbox}. + * + *

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. + * + *

{@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) { +} diff --git a/bridged/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java b/bridged/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java index d1cd9dc..53052cc 100644 --- a/bridged/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java +++ b/bridged/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java @@ -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"); diff --git a/bridged/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java b/bridged/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java new file mode 100644 index 0000000..bc015ac --- /dev/null +++ b/bridged/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java @@ -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}. + * + *

Proves the mechanism this ticket adds: a {@link LeadMessage} published to + * {@code lead..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 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 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 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 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..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 awaitPeek(LeadMailbox inbox) throws InterruptedException { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10); + List msgs = inbox.peek(); + while (msgs.isEmpty() && System.nanoTime() < deadline) { + Thread.sleep(50); + msgs = inbox.peek(); + } + return msgs; + } +}