broker: uriEnv config + non-fatal unreachable broker on boot
CI / build (pull_request) Failing after 56s
CI / contract (pull_request) Successful in 1m6s

#151: Broker gains uriEnv beside uri, taking the AMQP URI from an env var so
the password stays out of fleetd.yaml (same pattern as auth.tokenEnv). uriEnv
wins when set; isConfigured() treats a uriEnv naming an unset/blank variable
as unconfigured. A configured uriEnv is added to the startup required-secrets
report. Never logs the resolved URI (it carries the password).

#152: AmqpReplyInbox.open throwing at boot no longer stops the daemon. The
selection at the call site catches the failure and falls back to the in-memory
inbox for the process lifetime, warning loudly that durable cross-restart
delivery is off and logging the failed URI with credentials stripped.
This commit is contained in:
Dai Ha
2026-08-23 13:49:26 +02:00
parent 83b50753fe
commit 95a8dbcea9
4 changed files with 335 additions and 19 deletions
+6 -1
View File
@@ -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,
@@ -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 <em>loudly</em> — 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<String, String> 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 <em>with credentials stripped</em>. Never retries in the background — a broker that
* drops <em>after</em> 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 <em>below</em> 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 <em>below</em> 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;
}
@@ -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 <em>name</em> 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
* <em>not</em> 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<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 AmqpReplyInbox#DEFAULT_PREFETCH} when unset. */
@@ -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<ILoggingEvent> 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<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
return appender;
}
private static void detach(ListAppender<ILoggingEvent> appender) {
((Logger) LoggerFactory.getLogger(Fleetd.class)).detachAppender(appender);
}
private static void assertNoLogContains(ListAppender<ILoggingEvent> 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<ILoggingEvent> 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<ILoggingEvent> 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<String, String> env, boolean unreachable, Check check) {
RecordingAmqp opener = new RecordingAmqp();
opener.unreachable = unreachable;
ListAppender<ILoggingEvent> 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<ILoggingEvent> appender, ReplyInbox inbox);
}
}