diff --git a/bridged/fleetd.example.yaml b/bridged/fleetd.example.yaml
index 412390e..eea90c6 100644
--- a/bridged/fleetd.example.yaml
+++ b/bridged/fleetd.example.yaml
@@ -624,11 +624,16 @@ guard:
# RabbitMQ speaks the same AMQP 0-9-1, so it is a URI-only swap.
# uri → AMQP connection URI. No trailing slash ⇒ the default vhost "/"; an empty path ("/")
# is vhost "" and will NOT connect. Encode a named vhost as .../%2Fmyvhost.
+# uriEnv → CB-151: name of a host env var holding the AMQP URI, preferred over `uri` (wins
+# whenever set). The URI carries `user:pass@` inline, so naming a variable keeps the
+# password out of fleetd.yaml — same pattern as auth.tokenEnv/Profile.tokenEnv. A
+# uriEnv that resolves to an unset or blank variable is treated as NOT configured and
+# the daemon falls back to the in-memory inbox, warning loudly.
# prefetch → CB-527: consumer basicQos, capping how many unacked messages the inbox holds
# in-heap per owned target (the rest sits on the broker's durable queue instead of
# growing the JVM heap). Default 32 when omitted.
# broker:
-# uri: amqp://guest:guest@127.0.0.1:5672
+# uriEnv: LAVINMQ_URI
# prefetch: 32
# Active push-to-primary (CB-307 Stage 3). When a worker reply lands with no open fleet_send,
diff --git a/bridged/src/main/java/dev/ltms/fleet/Fleetd.java b/bridged/src/main/java/dev/ltms/fleet/Fleetd.java
index 72c118c..65dd744 100644
--- a/bridged/src/main/java/dev/ltms/fleet/Fleetd.java
+++ b/bridged/src/main/java/dev/ltms/fleet/Fleetd.java
@@ -371,18 +371,10 @@ public final class Fleetd {
StatusPoller poller = new StatusPoller(agents, injector, Injector.POLL_INTERVAL_MILLIS);
poller.start();
- // CB-307: reply inbox. A broker: block (with a uri) selects the AMQP-backed durable adapter;
- // absent, bridged stays soft-state on the in-memory inbox. The AMQP inbox owns a broker
+ // CB-307: reply inbox. A broker: block selects the AMQP-backed durable adapter; absent (or
+ // unusable), bridged stays soft-state on the in-memory inbox. The AMQP inbox owns a broker
// connection, so keep the reference to close it in the ordered shutdown hook.
- final ReplyInbox replyInbox;
- if (cfg.broker() != null && cfg.broker().isConfigured()) {
- replyInbox = AmqpReplyInbox.open(cfg.broker().uri(), cfg.broker().prefetchOrDefault());
- log.info("reply inbox: AMQP broker (durable) at {} (prefetch={})",
- cfg.broker().uri(), cfg.broker().prefetchOrDefault());
- } else {
- replyInbox = new InMemoryReplyInbox();
- log.info("reply inbox: in-memory (soft-state)");
- }
+ final ReplyInbox replyInbox = selectReplyInbox(cfg.broker(), System.getenv(), AmqpReplyInbox::open);
// CB-307: learn the primary's terminal from orchestration tool calls (or pin from config).
// The pin also feeds CallerResolver below: a primary running inside a herdr pane would
// otherwise resolve as a worker and be refused every orchestration tool.
@@ -577,14 +569,101 @@ public final class Fleetd {
return target -> presence.isPresent(target) || leads.get().containsKey(target);
}
+ /** Injection seam for {@link #selectReplyInbox}: production binds {@link AmqpReplyInbox#open}. */
+ @FunctionalInterface
+ interface AmqpOpener {
+ ReplyInbox open(String uri, int prefetch);
+ }
+
+ /**
+ * CB-151/152: pick the reply inbox. A usable broker — a literal {@code uri}, or a {@code
+ * uriEnv} whose variable resolves (both read from {@code env}) — selects the durable AMQP inbox.
+ * Everything else falls back to the in-memory inbox: no broker block, a blank {@code uri}, a
+ * {@code uriEnv} whose variable is unset or blank, or a broker unreachable at boot. The two
+ * lossy paths warn loudly — never silently — because what is lost is durable,
+ * cross-restart reply delivery. Package-private and env-injected so the selection is testable
+ * without a real broker or a mutable process environment.
+ */
+ static ReplyInbox selectReplyInbox(FleetConfig.Broker broker, Map env, AmqpOpener amqp) {
+ if (broker == null) {
+ log.info("reply inbox: in-memory (soft-state)");
+ return new InMemoryReplyInbox();
+ }
+ if (broker.hasUriEnv()) {
+ // uriEnv is authoritative whenever set (CB-151): the operator moved off clear text, so
+ // it must not quietly fall back onto a stale literal uri.
+ if (broker.uri() != null && !broker.uri().isBlank()) {
+ log.info("broker.uri is ignored because broker.uriEnv={} is set", broker.uriEnv());
+ }
+ String effectiveUri = broker.effectiveUri(env);
+ if (effectiveUri == null) {
+ log.warn("broker.uriEnv={} is unset or blank — durable AMQP reply inbox DISABLED. "
+ + "Replies are soft-state and will not survive a restart. Set {} in the "
+ + "daemon's environment (see scripts/redeploy-bridged.sh) and restart to "
+ + "use the durable broker inbox.",
+ broker.uriEnv(), broker.uriEnv());
+ log.info("reply inbox: in-memory (soft-state)");
+ return new InMemoryReplyInbox();
+ }
+ log.info("reply inbox: AMQP broker (durable) via env var {} (prefetch={})",
+ broker.uriEnv(), broker.prefetchOrDefault());
+ return openAmqpOrFallback(effectiveUri, broker.prefetchOrDefault(), "uriEnv " + broker.uriEnv(), amqp);
+ }
+ // No uriEnv: the literal uri path (existing behaviour).
+ if (broker.effectiveUri(env) == null) {
+ log.info("reply inbox: in-memory (soft-state)");
+ return new InMemoryReplyInbox();
+ }
+ log.info("reply inbox: AMQP broker (durable) (prefetch={})", broker.prefetchOrDefault());
+ return openAmqpOrFallback(broker.uri(), broker.prefetchOrDefault(), "uri", amqp);
+ }
+
+ /**
+ * Open the AMQP inbox, falling back to the in-memory inbox for this process lifetime if the
+ * broker cannot be reached at boot (CB-152). Not silent: the warning says durable delivery is
+ * off, replies are soft-state and will not survive a restart, plus the source that failed and
+ * the URI with credentials stripped. Never retries in the background — a broker that
+ * drops after startup already self-heals via the connection factory's automatic
+ * recovery; only the boot path is changed here.
+ */
+ private static ReplyInbox openAmqpOrFallback(String effectiveUri, int prefetch, String source,
+ AmqpOpener amqp) {
+ try {
+ return amqp.open(effectiveUri, prefetch);
+ } catch (IllegalStateException e) {
+ log.warn("cannot reach AMQP broker ({}, {}) — falling back to the in-memory reply inbox "
+ + "for this process lifetime. Durable, cross-restart reply delivery is OFF; "
+ + "replies are soft-state and will not survive a restart. Reason: {}",
+ source, stripCredentials(effectiveUri), reasonOf(e));
+ log.info("reply inbox: in-memory (soft-state)");
+ return new InMemoryReplyInbox();
+ }
+ }
+
+ /** An AMQP URI carries {@code user:pass@} inline — show the host/port, never the credentials. */
+ static String stripCredentials(String uri) {
+ return uri == null ? null : uri.replaceAll("://[^@/]*@", "://");
+ }
+
+ /** The deepest cause's class and message — the outermost {@code IllegalStateException} echoes the URI (with password). */
+ private static String reasonOf(Throwable e) {
+ Throwable t = e;
+ while (t.getCause() != null && t.getCause() != t) {
+ t = t.getCause();
+ }
+ String msg = t.getMessage();
+ return t.getClass().getSimpleName() + (msg == null || msg.isBlank() ? "" : ": " + msg);
+ }
+
/**
* CB-594: which env vars the loaded config actually needs, and why — every non-{@code
* subscription} profile's {@code tokenEnv} (a subscription profile never reads one, see
* {@link FleetConfig.Profile#isSubscription()}), plus every profile's {@code gitTokenEnv}
- * where set (opt-in). Derived from the config, not hard-coded, so a new profile is covered for
- * free. A var required by more than one profile is one entry naming every profile that needs
- * it. Deliberately excludes {@code auth.tokenEnv}: that one is already enforced loudly, by a
- * startup throw in {@code main()} — about 370 lines below this method's call site
+ * where set (opt-in), plus a configured {@code broker.uriEnv} (CB-151). Derived from the
+ * config, not hard-coded, so a new profile is covered for free. A var required by more than one
+ * profile is one entry naming every profile that needs it. Deliberately excludes {@code
+ * auth.tokenEnv}: that one is already enforced loudly, by a startup throw in {@code main()} —
+ * about 370 lines below this method's call site
* ({@link #reportRequiredSecrets(FleetConfig)}), not a few lines above it. That throw only
* fires when {@code auth.mode: token} is configured; under the default loopback-trust mode it
* never runs, and {@code auth.tokenEnv} is simply not required.
@@ -604,6 +683,11 @@ public final class Fleetd {
.add("profile '" + name + "' gitTokenEnv");
}
});
+ FleetConfig.Broker broker = cfg.broker();
+ if (broker != null && broker.hasUriEnv()) {
+ requiredBy.computeIfAbsent(broker.uriEnv(), _ -> new ArrayList<>())
+ .add("broker uriEnv");
+ }
return requiredBy;
}
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 33ad132..ea84552 100644
--- a/bridged/src/main/java/dev/ltms/fleet/config/FleetConfig.java
+++ b/bridged/src/main/java/dev/ltms/fleet/config/FleetConfig.java
@@ -551,16 +551,52 @@ public record FleetConfig(
*
* @param uri AMQP connection URI, e.g. {@code amqp://guest:guest@127.0.0.1:5672/}. Blank/
* {@code null} ⇒ the broker block is treated as absent (in-memory adapter).
+ * Ignored when {@code uriEnv} is set.
+ * @param uriEnv name of a host env var holding the AMQP URI (CB-151). The URI carries its
+ * credentials inline, so giving the variable name keeps the password
+ * out of the config file, same as {@code auth.tokenEnv}/{@code Profile.tokenEnv}.
+ * Wins over {@code uri} whenever set. Blank/{@code null} ⇒ ignored.
* @param prefetch CB-527: the consumer's {@code basicQos} prefetch count, bounding how many
* unacked messages the AMQP inbox holds in-heap per owned target. {@code null}/
* non-positive ⇒ {@link AmqpReplyInbox#DEFAULT_PREFETCH}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
- public record Broker(String uri, Integer prefetch) {
+ public record Broker(String uri, String uriEnv, Integer prefetch) {
- /** True when a usable broker URI is configured (an empty block does not enable AMQP). */
+ /** 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 broker URI is configured (an empty block does not enable AMQP).
+ * Honors {@code uriEnv} first: if it names a variable that is unset or blank, the broker is
+ * not configured (the daemon falls back to the in-memory inbox) — a bare {@code uri}
+ * is only consulted when no {@code uriEnv} is set. Env lookup makes this process-dependent;
+ * callers already reading {@link System#getenv} are the right ones to invoke it.
+ */
public boolean isConfigured() {
- return uri != null && !uri.isBlank();
+ 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 AmqpReplyInbox#DEFAULT_PREFETCH} when unset. */
diff --git a/bridged/src/test/java/dev/ltms/fleet/FleetdReplyInboxSelectionTest.java b/bridged/src/test/java/dev/ltms/fleet/FleetdReplyInboxSelectionTest.java
new file mode 100644
index 0000000..f34adb6
--- /dev/null
+++ b/bridged/src/test/java/dev/ltms/fleet/FleetdReplyInboxSelectionTest.java
@@ -0,0 +1,191 @@
+package dev.ltms.fleet;
+
+import ch.qos.logback.classic.Level;
+import ch.qos.logback.classic.Logger;
+import ch.qos.logback.classic.spi.ILoggingEvent;
+import ch.qos.logback.core.read.ListAppender;
+import dev.ltms.fleet.config.FleetConfig;
+import dev.ltms.fleet.msg.AmqpReplyInbox;
+import dev.ltms.fleet.msg.InMemoryReplyInbox;
+import dev.ltms.fleet.msg.ReplyInbox;
+import org.junit.jupiter.api.Test;
+import org.slf4j.LoggerFactory;
+
+import java.net.ServerSocket;
+import java.util.List;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * CB-151/152: the reply inbox is selected in {@link Fleetd#selectReplyInbox}, NOT in the config
+ * record — a test that only exercises {@code broker.isConfigured()} would pass even if {@code
+ * Fleetd} never honoured {@code uriEnv}. These tests drive the real selection logic with an
+ * injected env map and an injected AMQP opener, so they prove which inbox the daemon actually
+ * picks, and that it never logs the resolved URI (which carries the password).
+ */
+class FleetdReplyInboxSelectionTest {
+
+ public static final String SECRET = "s3cr3tPw";
+ /** A resolved URI whose userinfo carries the password, so we can assert it never leaks. */
+ private static final String RESOLVED_URI = "amqp://user:" + SECRET + "@broker.example:5672/vh";
+
+ /** Fake AMQP opener: records the URI it was offered, or fails as if the broker were unreachable. */
+ private static final class RecordingAmqp implements Fleetd.AmqpOpener {
+ String offeredUri;
+ boolean unreachable;
+ final ReplyInbox inbox = new InMemoryReplyInbox();
+
+ @Override
+ public ReplyInbox open(String uri, int prefetch) {
+ if (unreachable) {
+ throw new IllegalStateException("cannot connect to AMQP broker at " + uri,
+ new java.net.ConnectException("Connection refused"));
+ }
+ this.offeredUri = uri;
+ return inbox;
+ }
+ }
+
+ private static ListAppender attach() {
+ Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
+ // logback-test.xml pins dev.ltms.fleet to WARN; raise it so INFO selection lines are captured.
+ logger.setLevel(Level.INFO);
+ ListAppender appender = new ListAppender<>();
+ appender.start();
+ logger.addAppender(appender);
+ return appender;
+ }
+
+ private static void detach(ListAppender appender) {
+ ((Logger) LoggerFactory.getLogger(Fleetd.class)).detachAppender(appender);
+ }
+
+ private static void assertNoLogContains(ListAppender appender, String secret) {
+ assertTrue(appender.list.stream().noneMatch(e -> e.getFormattedMessage().contains(secret)),
+ "no log line may contain the resolved URI's password");
+ }
+
+ @Test
+ void uriEnvSetAndPresentSelectsAmqpWithTheResolvedUri() {
+ FleetConfig.Broker broker = new FleetConfig.Broker(null, "LAVINMQ_URI", null);
+ recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), false, (opener, appender, inbox) -> {
+ assertEquals(RESOLVED_URI, opener.offeredUri,
+ "the daemon must connect with the value resolved from uriEnv — selection, not just parse");
+ assertEquals(opener.inbox, inbox, "the AMQP opener's inbox is what is selected");
+ assertNoLogContains(appender, SECRET);
+ });
+ }
+
+ @Test
+ void uriEnvSetButVariableMissingFallsBackToInMemoryAndWarns() {
+ FleetConfig.Broker broker = new FleetConfig.Broker(null, "LAVINMQ_URI", null);
+ recording(broker, Map.of(), false, (opener, appender, inbox) -> {
+ assertInstanceOf(InMemoryReplyInbox.class, inbox);
+ assertNull(opener.offeredUri, "AMQP must never be attempted when the variable is missing");
+ assertTrue(hasWarnContaining(appender, "LAVINMQ_URI") && hasWarnContaining(appender, "DISABLED"),
+ "a missing uriEnv variable must warn loudly, not fail silently");
+ });
+ }
+
+ @Test
+ void uriEnvSetButVariableBlankFallsBackToInMemoryAndWarns() {
+ FleetConfig.Broker broker = new FleetConfig.Broker("amqp://user:lame@old:5672/", "LAVINMQ_URI", null);
+ recording(broker, Map.of("LAVINMQ_URI", " "), false, (opener, appender, inbox) -> {
+ assertInstanceOf(InMemoryReplyInbox.class, inbox);
+ assertNull(opener.offeredUri, "a blank env value must not select AMQP, not even via the literal uri");
+ assertTrue(hasWarnContaining(appender, "LAVINMQ_URI"),
+ "a blank uriEnv value must warn, and must not fall back to the literal uri");
+ });
+ }
+
+ @Test
+ void bothUriAndUriEnvSetUriEnvWinsDeterministically() {
+ FleetConfig.Broker broker
+ = new FleetConfig.Broker("amqp://user:oldpw@old.example:5672/", "LAVINMQ_URI", null);
+ recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), false, (opener, appender, inbox) -> {
+ assertEquals(RESOLVED_URI, opener.offeredUri,
+ "uriEnv must win over uri, deterministically, every run");
+ assertTrue(appender.list.stream().anyMatch(e -> e.getFormattedMessage().contains("broker.uri is ignored")),
+ "must log that the literal uri is ignored when uriEnv is set");
+ assertNoLogContains(appender, SECRET);
+ assertNoLogContains(appender, "oldpw");
+ });
+ }
+
+ @Test
+ void unreachableBrokerStartsDaemonWithInMemoryInboxAndLoudWarning() {
+ FleetConfig.Broker broker = new FleetConfig.Broker(null, "LAVINMQ_URI", null);
+ recording(broker, Map.of("LAVINMQ_URI", RESOLVED_URI), true, (opener, appender, inbox) -> {
+ assertInstanceOf(InMemoryReplyInbox.class, inbox, "an unreachable broker must NOT stop the daemon");
+ String warn = appender.list.stream()
+ .filter(e -> e.getLevel() == Level.WARN)
+ .map(ILoggingEvent::getFormattedMessage)
+ .reduce("", (a, b) -> a + "\n" + b)
+ .toLowerCase();
+ assertTrue(warn.contains("durable") && warn.contains("soft-state"),
+ "the warning must say exactly what was lost: durable delivery off, replies soft-state");
+ assertTrue(!warn.contains(SECRET), "the failing URI must be logged with credentials stripped");
+ assertNoLogContains(appender, SECRET);
+ });
+ }
+
+ @Test
+ void noBrokerConfiguredStaysQuietInMemory() {
+ FleetConfig.Broker broker = null;
+ recording(broker, Map.of(), false, (opener, appender, inbox) -> {
+ assertInstanceOf(InMemoryReplyInbox.class, inbox);
+ assertTrue(appender.list.stream().noneMatch(e -> e.getLevel() == Level.WARN),
+ "no broker configured must keep the existing QUIET in-memory path — no warning");
+ });
+ }
+
+ @Test
+ void aRealUnreachableBrokerFallsBackViaTheRealOpener() throws Exception {
+ // A guaranteed-closed port: grab an ephemeral one, release it, then connect to the now-dead
+ // address. This exercises AmqpReplyInbox.open's real throw path without any container.
+ int closedPort;
+ try (ServerSocket s = new ServerSocket(0)) {
+ closedPort = s.getLocalPort();
+ }
+ FleetConfig.Broker broker = new FleetConfig.Broker(null, "LAVINMQ_URI", null);
+ ListAppender appender = attach();
+ try {
+ ReplyInbox inbox = Fleetd.selectReplyInbox(
+ broker, Map.of("LAVINMQ_URI", "amqp://user:" + SECRET + "@127.0.0.1:" + closedPort + "/vh"),
+ AmqpReplyInbox::open);
+ assertInstanceOf(InMemoryReplyInbox.class, inbox,
+ "a genuinely unreachable broker (real AmqpReplyInbox::open) must fall back to in-memory");
+ } finally {
+ detach(appender);
+ }
+ assertNoLogContains(appender, SECRET);
+ }
+
+ private boolean hasWarnContaining(ListAppender appender, String fragment) {
+ return appender.list.stream().anyMatch(e ->
+ e.getLevel() == Level.WARN && e.getFormattedMessage().contains(fragment));
+ }
+
+ /** Runs one selection under a captured log, asserting on its outcome. */
+ private void recording(FleetConfig.Broker broker, Map env, boolean unreachable, Check check) {
+ RecordingAmqp opener = new RecordingAmqp();
+ opener.unreachable = unreachable;
+ ListAppender appender = attach();
+ ReplyInbox inbox;
+ try {
+ inbox = Fleetd.selectReplyInbox(broker, env, opener);
+ } finally {
+ detach(appender);
+ }
+ check.run(opener, appender, inbox);
+ }
+
+ @FunctionalInterface
+ private interface Check {
+ void run(RecordingAmqp opener, ListAppender appender, ReplyInbox inbox);
+ }
+}