From 95a8dbcea9a4e9b59e457c79a2f54b51268260c6 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sun, 23 Aug 2026 13:49:26 +0200 Subject: [PATCH] broker: uriEnv config + non-fatal unreachable broker on boot #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. --- bridged/fleetd.example.yaml | 7 +- .../src/main/java/dev/ltms/fleet/Fleetd.java | 114 +++++++++-- .../dev/ltms/fleet/config/FleetConfig.java | 42 +++- .../fleet/FleetdReplyInboxSelectionTest.java | 191 ++++++++++++++++++ 4 files changed, 335 insertions(+), 19 deletions(-) create mode 100644 bridged/src/test/java/dev/ltms/fleet/FleetdReplyInboxSelectionTest.java 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); + } +}