Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 3ba6d6784c CB-594: make supervision and a working fleet possible at the same time
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 1m40s
Adds scripts/bridged-launchd-wrapper.sh so the launchd-run daemon still gets
WORKER_GITEA_TOKEN/AI_GATEWAY_TOKEN by execing through a login shell (launchd
never sources secrets.sh itself). bridged now logs at startup which required
token env vars (derived from each profile's tokenEnv/gitTokenEnv, not a
hand-written list) resolved or are MISSING, by name only. Fills in the real
paths in deploy/dev.ltms.bridged.plist for this host and points it at the
wrapper. scripts/redeploy-bridged.sh now detects a loaded launchd agent and
uses launchctl unload/load instead of a raw kill+nohup, because a bare
SIGTERM exits this JVM at 143 (measured) which KeepAlive.SuccessfulExit=false
reads as a crash and would race the script's own restart; --check reports
installed/loaded state and stays read-only.
2026-08-16 17:13:17 +02:00
9 changed files with 338 additions and 304 deletions
+2 -6
View File
@@ -434,14 +434,10 @@ guard:
# on a durable per-target queue (agent.<target>.inbox) and survive a restart — the broker
# redelivers anything the primary had not yet drained. Production default is LavinMQ; a stock
# 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.
# 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.
# 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.
# broker:
# uri: amqp://guest:guest@127.0.0.1:5672
# prefetch: 32
# Active push-to-primary (CB-307 Stage 3). When a worker reply lands with no open bridge_send,
# the ReplyPushLoop injects a *drain nudge* (never the payload) into the primary's own herdr
@@ -83,6 +83,10 @@ public final class Bridged {
static void main(String[] args) {
Path configPath = Path.of(args.length > 0 ? args[0] : "bridged.yaml");
BridgedConfig cfg = BridgedConfig.load(configPath);
// CB-594: report which secret env vars the config actually needs, by name, before anything
// else can fail on a silently-empty one. A daemon started without a login shell (launchd)
// boots fine either way — this is the only thing that says so out loud.
reportRequiredSecrets(cfg);
// CB-559: `cfg` stays the startup snapshot — every validation and every piece of one-time
// wiring below reads it, and must, because those decisions cannot be unmade. `config` is the
// live reference the hot paths read per use. Which keys can actually move is ConfigRef's
@@ -348,9 +352,8 @@ public final class Bridged {
// 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());
replyInbox = AmqpReplyInbox.open(cfg.broker().uri());
log.info("reply inbox: AMQP broker (durable) at {}", cfg.broker().uri());
} else {
replyInbox = new InMemoryReplyInbox();
log.info("reply inbox: in-memory (soft-state)");
@@ -544,6 +547,63 @@ public final class Bridged {
return target -> presence.isPresent(target) || leads.get().containsKey(target);
}
/**
* 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 BridgedConfig.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, a few lines above this method's call site.
*
* <p>Package-private and pure (no I/O, no logging) so the derivation is unit-testable without
* capturing log output; {@link #reportRequiredSecrets(BridgedConfig)} is the logging caller.
*/
static Map<String, List<String>> requiredSecretEnvVars(BridgedConfig cfg) {
Map<String, List<String>> requiredBy = new LinkedHashMap<>();
cfg.profiles().forEach((name, profile) -> {
if (!profile.isSubscription()) {
requiredBy.computeIfAbsent(profile.tokenEnv(), _ -> new ArrayList<>())
.add("profile '" + name + "' tokenEnv");
}
if (profile.hasGitToken()) {
requiredBy.computeIfAbsent(profile.gitTokenEnv(), _ -> new ArrayList<>())
.add("profile '" + name + "' gitTokenEnv");
}
});
return requiredBy;
}
/**
* CB-594: log, by name only, which required env vars (see {@link #requiredSecretEnvVars}) are
* set in the daemon's own process environment — the environment every profile's {@code
* tokenEnv}/{@code gitTokenEnv} is read from at spawn time (see
* {@code HerdrPeerLauncher.resolveEnv}). Never logs a value, a prefix, or a length.
*
* <p>A missing entry only warns — it must never refuse to start. A daemon that boots and says
* what is wrong is strictly more useful than one that will not boot at all.
*/
private static void reportRequiredSecrets(BridgedConfig cfg) {
Map<String, List<String>> requiredBy = requiredSecretEnvVars(cfg);
if (requiredBy.isEmpty()) {
log.info("startup secrets: no profile references a token env var — nothing to check");
return;
}
Map<String, String> env = System.getenv();
requiredBy.forEach((varName, sources) -> {
String value = env.get(varName);
if (value != null && !value.isBlank()) {
log.info("startup secret {}: set ({})", varName, String.join(", ", sources));
} else {
log.warn("startup secret {}: MISSING ({}) — the daemon will start anyway, and this "
+ "failure stays invisible until a worker actually needs it. Fix "
+ "${SHARED_ENV}/tools/secrets.sh and restart bridged from a LOGIN "
+ "shell (see scripts/redeploy-bridged.sh).",
varName, String.join(", ", sources));
}
});
}
/**
* Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504).
*
@@ -5,7 +5,6 @@ import com.fasterxml.jackson.core.JsonParser;
import com.fasterxml.jackson.core.JsonToken;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
import dev.ltms.bridged.msg.AmqpReplyInbox;
import dev.ltms.bridged.peer.MemberRole;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -533,24 +532,16 @@ public record BridgedConfig(
* stays soft-state. Production default is LavinMQ; a stock RabbitMQ speaks the same AMQP 0-9-1
* and is a URI-only swap.
*
* @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).
* @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}.
* @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).
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Broker(String uri, Integer prefetch) {
public record Broker(String uri) {
/** True when a usable broker URI is configured (an empty block does not enable AMQP). */
public boolean isConfigured() {
return uri != null && !uri.isBlank();
}
/** The prefetch to use, defaulting to {@link AmqpReplyInbox#DEFAULT_PREFETCH} when unset. */
public int prefetchOrDefault() {
return (prefetch != null && prefetch > 0) ? prefetch : AmqpReplyInbox.DEFAULT_PREFETCH;
}
}
/**
@@ -7,7 +7,6 @@ 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;
@@ -15,13 +14,7 @@ import java.io.IOException;
import java.nio.charset.StandardCharsets;
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;
/**
* AMQP-backed {@link ReplyInbox} (CB-307 Stage 2): genuine cross-restart durability behind the same
@@ -36,27 +29,6 @@ import java.util.concurrent.TimeoutException;
* leaves them on the broker — it redelivers on reconnect. That is the durability the in-memory
* adapter cannot give, with the port contract preserved.
*
* <p><strong>Prefetch bounds the held backlog (CB-527).</strong> The consumer channel calls
* {@code basicQos} with a configurable prefetch count ({@link #DEFAULT_PREFETCH} unless the caller
* passes another value to {@link #open(String, int)}) before starting any consumer. Without a bound,
* the broker pushes its entire queue into {@link #held} the instant a target is {@link #own owned},
* so an undrained primary grows the JVM heap without limit and any queue-level control
* ({@code x-max-length}, per-message TTL) never fires because the queue never actually holds a
* backlog. Prefetch keeps the backlog where it is visible — on the broker — until the owner drains it.
*
* <p><strong>Publishes require a confirmed, routable delivery (CB-528).</strong> {@link #publish}
* runs on a channel separate from the consume/ack channel ({@link #channel}), so a slow or blocked
* publish confirm can never hold {@link #channelLock} and stall an ack — the ack path never waits on
* a publish confirm. That publish channel is in publisher-confirm mode and every publish sets the
* {@code mandatory} flag, so an unroutable publish (queue not declared, e.g. the owner never called
* {@link #own}) is returned by the broker instead of silently dropped. The broker sends the
* <em>return</em> for an unroutable message before the <em>confirm</em> that covers it — the ack/nack
* callback checks the returned-set at confirm time rather than assuming an ack means routed — so
* "confirmed" here means "durably queued", not merely "accepted by the broker". A returned or nacked
* (or un-confirmed within the timeout) publish surfaces as an {@link IllegalStateException} on the
* caller's thread; the caller — {@link MessageService#reply} — must not report success for a
* black-holed reply.
*
* <p><strong>Ownership is explicit.</strong> {@link #own} declares the queue and starts the consumer;
* {@link #release} cancels it. {@link #publish} sends to the queue but does <em>not</em> imply ownership
* and does not attach a consumer. This split is required by CB-308 federation, where one gateway may
@@ -82,12 +54,6 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
private static final String QUEUE_PREFIX = "agent.";
private static final String QUEUE_SUFFIX = ".inbox";
/** CB-527: the prefetch used when a caller does not pass an explicit value to {@link #open(String, int)}. */
public static final int DEFAULT_PREFETCH = 32;
/** How long {@link #publish} waits for its publisher confirm before failing the call (CB-528). */
private static final long CONFIRM_TIMEOUT_MS = 10_000L;
private final Connection connection;
private final Channel channel;
/** All channel operations (publish/declare/ack/cancel) serialize on this — a Channel is not thread-safe. */
@@ -97,81 +63,39 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
/** Targets whose queue is declared and consumer is running, mapped to their broker consumer tag. */
private final ConcurrentHashMap<String, String> consumerTags = new ConcurrentHashMap<>();
/**
* CB-528: a dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume
* + ack) so a publish confirm round trip never blocks under {@link #channelLock} and stalls an ack.
*/
private final Channel publishChannel;
private final Object publishChannelLock = new Object();
/** In-flight publishes awaiting their confirm, keyed by the publish channel's sequence number. */
private final ConcurrentSkipListMap<Long, Pending> pendingBySeq = new ConcurrentSkipListMap<>();
/** The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery tag. */
private final ConcurrentHashMap<String, Pending> pendingByMsgId = new ConcurrentHashMap<>();
/** A message pulled off the broker but not yet acked: its delivery-tag plus the port payload. */
private record Held(long deliveryTag, InboxMessage message) {}
/** A publish awaiting its confirm; {@link #returned} records whether the broker already returned it. */
private static final class Pending {
final String msgId;
final CompletableFuture<Void> confirmed = new CompletableFuture<>();
volatile boolean returned;
Pending(String msgId) {
this.msgId = msgId;
}
}
/** Connect to {@code uri} (e.g. {@code amqp://guest:guest@127.0.0.1:5672/}) with {@link #DEFAULT_PREFETCH}. */
/** Connect to {@code uri} (e.g. {@code amqp://guest:guest@127.0.0.1:5672/}) and open the inbox. */
public static AmqpReplyInbox open(String uri) {
return open(uri, DEFAULT_PREFETCH);
}
/** As {@link #open(String)}, with an explicit consumer prefetch (CB-527: caps the held backlog per target). */
public static AmqpReplyInbox open(String uri, int prefetch) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
return new AmqpReplyInbox(factory.newConnection("bridged-reply-inbox"), prefetch);
return new AmqpReplyInbox(factory.newConnection("bridged-reply-inbox"));
} catch (Exception e) {
throw new IllegalStateException("cannot connect to AMQP broker at " + uri, e);
}
}
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for the contract test). */
/** Wrap an already-open connection (injection seam for the contract test). */
AmqpReplyInbox(Connection connection) {
this(connection, DEFAULT_PREFETCH);
}
/** As above, with an explicit prefetch (injection seam for the contract test). */
AmqpReplyInbox(Connection connection, int prefetch) {
this.connection = connection;
try {
this.channel = connection.createChannel();
// CB-527: bound the held backlog per owned target — must be set before any own()/basicConsume.
this.channel.basicQos(prefetch);
this.publishChannel = connection.createChannel();
this.publishChannel.confirmSelect();
this.publishChannel.addReturnListener(this::onReturn);
this.publishChannel.addConfirmListener(this::onAck, this::onNack);
} 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 — its sequence number
// meant nothing on the old channel and means nothing on the recovered one, so fail it now
// rather than let it silently ride out CONFIRM_TIMEOUT_MS.
// repopulates it with valid tags (dedup by msgId still prevents any double-queue).
if (connection instanceof Recoverable recoverable) {
recoverable.addRecoveryListener(new RecoveryListener() {
@Override
public void handleRecovery(Recoverable recoverable) {
held.clear();
failPendingPublishesOnRecovery();
log.info("AMQP connection recovered; cleared held replies for fresh redelivery");
}
@@ -217,12 +141,6 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
}
/**
* Publish {@code content} and block until the broker's publisher confirm for it lands (CB-528).
* 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.
*/
@Override
public void publish(String target, String msgId, String content) {
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
@@ -230,35 +148,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
.deliveryMode(2) // persistent — survives a broker restart
.contentType("text/plain")
.build();
Pending pending = new Pending(msgId);
long seq;
synchronized (publishChannelLock) {
seq = publishChannel.getNextPublishSeqNo();
pendingBySeq.put(seq, pending);
pendingByMsgId.put(msgId, pending);
try {
publishChannel.basicPublish("", queueName(target), true, props,
content.getBytes(StandardCharsets.UTF_8));
} catch (IOException e) {
pendingBySeq.remove(seq, pending);
pendingByMsgId.remove(msgId, pending);
throw new IllegalStateException("cannot publish reply to " + queueName(target), 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 reply " + msgId + " to " + queueName(target)
+ " 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 " + msgId, e);
} finally {
pendingBySeq.remove(seq, pending);
pendingByMsgId.remove(msgId, pending);
synchronized (channelLock) {
channel.basicPublish("", queueName(target), props, content.getBytes(StandardCharsets.UTF_8));
}
} catch (IOException e) {
throw new IllegalStateException("cannot publish reply to " + queueName(target), e);
}
}
@@ -327,65 +222,6 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
};
}
/** Broker return for an unroutable {@code mandatory} publish — arrives BEFORE its confirm (CB-528). */
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 reply {} (routingKey={}, {} {}) with no matching in-flight publish"
+ " — already resolved by a prior confirm", msgId, r.getRoutingKey(), r.getReplyCode(),
r.getReplyText());
}
}
private void onAck(long seq, boolean multiple) {
resolveConfirm(seq, multiple, true);
}
private void onNack(long seq, boolean multiple) {
resolveConfirm(seq, multiple, false);
}
/**
* Resolve every pending publish covered by this confirm (a single seq, or — {@code multiple} —
* every seq up to and including it). Checks {@link Pending#returned} at confirm time: since the
* broker's return for an unroutable message always precedes its confirm, an ack that arrives after
* a return means "confirmed but never routed", not "durably queued".
*/
private void resolveConfirm(long seq, boolean multiple, boolean ack) {
NavigableMap<Long, Pending> covered = multiple
? pendingBySeq.headMap(seq, true)
: pendingBySeq.subMap(seq, true, seq, true);
for (var it = covered.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
if (ack && !pending.returned) {
pending.confirmed.complete(null);
} else if (ack) {
pending.confirmed.completeExceptionally(new IllegalStateException(
"reply " + pending.msgId + " was returned as unroutable (queue not declared/owned)"));
} else {
pending.confirmed.completeExceptionally(new IllegalStateException(
"broker nacked publish of reply " + pending.msgId));
}
}
}
/** Fail every publish still awaiting its confirm — their sequence numbers are stale after recovery. */
private void failPendingPublishesOnRecovery() {
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 reply " + pending.msgId
+ " is unknown"));
}
}
private static String queueName(String target) {
return QUEUE_PREFIX + target + QUEUE_SUFFIX;
}
@@ -397,11 +233,6 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
} catch (Exception e) {
log.debug("AMQP channel close: {}", e.toString());
}
try {
publishChannel.close();
} catch (Exception e) {
log.debug("AMQP publish channel close: {}", e.toString());
}
try {
connection.close();
} catch (Exception e) {
@@ -0,0 +1,114 @@
package dev.ltms.bridged;
import dev.ltms.bridged.config.BridgedConfig;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* CB-594: {@link Bridged#requiredSecretEnvVars(BridgedConfig)} is what decides what the startup
* secret report checks — it must derive that set from the config, not a hand-written list, or a
* new profile's token silently stops being reported.
*/
class RequiredSecretEnvVarsTest {
private static BridgedConfig load(Path dir, String yaml) throws Exception {
Path f = dir.resolve("bridged.yaml");
Files.writeString(f, yaml);
return BridgedConfig.load(f);
}
@Test
void collectsATokenEnvPerNonSubscriptionProfile(@TempDir Path dir) throws Exception {
BridgedConfig cfg = load(dir, """
profiles:
local:
baseUrl: http://gx00.gw:8000
tokenEnv: AI_GATEWAY_TOKEN
""");
Map<String, List<String>> required = Bridged.requiredSecretEnvVars(cfg);
assertTrue(required.containsKey("AI_GATEWAY_TOKEN"));
assertEquals(List.of("profile 'local' tokenEnv"), required.get("AI_GATEWAY_TOKEN"));
}
@Test
void aSubscriptionProfileNeedsNoTokenEnv(@TempDir Path dir) throws Exception {
BridgedConfig cfg = load(dir, """
profiles:
opus:
subscription: true
model: claude-opus-5
""");
assertTrue(Bridged.requiredSecretEnvVars(cfg).isEmpty(),
"subscription: true never reads ANTHROPIC_AUTH_TOKEN — see Profile#isSubscription");
}
@Test
void gitTokenEnvIsOptInAndCollectedWhenSet(@TempDir Path dir) throws Exception {
BridgedConfig cfg = load(dir, """
profiles:
local:
baseUrl: http://gx00.gw:8000
tokenEnv: AI_GATEWAY_TOKEN
gitTokenEnv: WORKER_GITEA_TOKEN
""");
Map<String, List<String>> required = Bridged.requiredSecretEnvVars(cfg);
assertTrue(required.containsKey("WORKER_GITEA_TOKEN"));
assertEquals(List.of("profile 'local' gitTokenEnv"), required.get("WORKER_GITEA_TOKEN"));
}
@Test
void noGitTokenEnvMeansNothingIsRequiredForIt(@TempDir Path dir) throws Exception {
BridgedConfig cfg = load(dir, """
profiles:
local:
baseUrl: http://gx00.gw:8000
tokenEnv: AI_GATEWAY_TOKEN
""");
assertFalse(Bridged.requiredSecretEnvVars(cfg).containsKey("WORKER_GITEA_TOKEN"));
}
@Test
void aVarSharedByTwoProfilesIsReportedOnceNamingBoth(@TempDir Path dir) throws Exception {
BridgedConfig cfg = load(dir, """
profiles:
local:
baseUrl: http://gx00.gw:8000
tokenEnv: AI_GATEWAY_TOKEN
gitTokenEnv: WORKER_GITEA_TOKEN
gx:
kind: opencode
baseUrl: https://llm.ltms.dev/v1
tokenEnv: AI_GATEWAY_TOKEN
gitTokenEnv: WORKER_GITEA_TOKEN
""");
Map<String, List<String>> required = Bridged.requiredSecretEnvVars(cfg);
assertEquals(List.of("profile 'local' tokenEnv", "profile 'gx' tokenEnv"),
required.get("AI_GATEWAY_TOKEN"));
assertEquals(List.of("profile 'local' gitTokenEnv", "profile 'gx' gitTokenEnv"),
required.get("WORKER_GITEA_TOKEN"));
}
@Test
void noProfilesMeansNothingIsRequired(@TempDir Path dir) throws Exception {
BridgedConfig cfg = load(dir, "bind:\n host: 127.0.0.1\n port: 8765\n");
assertTrue(Bridged.requiredSecretEnvVars(cfg).isEmpty());
}
}
@@ -16,7 +16,6 @@ import java.util.List;
import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
@@ -167,88 +166,6 @@ class AmqpReplyInboxContractTest {
}
}
@Test
void prefetchBoundsTheHeldBacklog() throws Exception {
String target = "worker-prefetch-" + System.nanoTime();
int prefetch = 4;
int published = 10;
try (AmqpReplyInbox inbox = new AmqpReplyInbox(newConnection(), prefetch);
Connection inspect = newConnection()) {
inbox.own(target);
for (int i = 0; i < published; i++) {
inbox.publish(target, "m" + i, "payload " + i);
}
awaitHeldAtLeast(inbox, target, prefetch);
long depth;
try (Channel ch = inspect.createChannel()) {
depth = ch.queueDeclarePassive(queueName(target)).getMessageCount();
}
assertTrue(depth >= published - prefetch,
"broker should still hold at least " + (published - prefetch)
+ " undelivered messages behind a prefetch of " + prefetch + ", saw " + depth);
drainUntilEmpty(inbox, target, inspect);
}
}
@Test
void unroutablePublishReportsFailureNotSilentSuccess() throws Exception {
String target = "worker-unroutable-" + System.nanoTime();
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
// Deliberately never own(target): the queue is never declared, so the default-exchange
// route to agent.<target>.inbox does not exist and the broker must return the publish.
IllegalStateException ex = assertThrows(IllegalStateException.class,
() -> inbox.publish(target, "m1", "nobody home"));
assertTrue(ex.getMessage() != null && ex.getMessage().toLowerCase().contains("unroutable"),
"expected an unroutable-publish failure, got: " + ex.getMessage());
}
}
@Test
void confirmedPublishDeliversNormally() throws Exception {
String target = "worker-confirm-" + System.nanoTime();
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
inbox.own(target);
inbox.publish(target, "m1", "confirmed delivery"); // must return normally: routed and confirmed
List<ReplyInbox.InboxMessage> got = awaitPeek(inbox, target);
assertEquals(1, got.size());
assertEquals("confirmed delivery", got.getFirst().content());
inbox.ack(target, "m1");
}
}
/** Poll peek until at least {@code n} replies for {@code target} are held, or ~10s elapse. */
@SuppressWarnings("BusyWait")
private static void awaitHeldAtLeast(AmqpReplyInbox inbox, String target, int n) throws InterruptedException {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
while (inbox.peek(target).size() < n && System.nanoTime() < deadline) {
Thread.sleep(50);
}
}
/**
* Repeatedly ack whatever is currently held (freeing prefetch slots for the next delivery) until
* both the local snapshot and the broker's own queue depth are empty, or ~10s elapse.
*/
@SuppressWarnings("BusyWait")
private static void drainUntilEmpty(AmqpReplyInbox inbox, String target, Connection inspect) throws Exception {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
while (System.nanoTime() < deadline) {
for (ReplyInbox.InboxMessage msg : inbox.peek(target)) {
inbox.ack(target, msg.msgId());
}
try (Channel ch = inspect.createChannel()) {
if (ch.queueDeclarePassive(queueName(target)).getMessageCount() == 0 && inbox.peek(target).isEmpty()) {
return;
}
}
Thread.sleep(50);
}
throw new AssertionError("did not drain " + target + " to empty within the deadline");
}
/** Poll peek (broker delivery is async) until a reply for {@code target} appears or ~10s elapse. */
@SuppressWarnings("BusyWait") // deliberate poll for async broker delivery, bounded by the deadline
private static List<ReplyInbox.InboxMessage> awaitPeek(AmqpReplyInbox inbox, String target)
+43 -13
View File
@@ -1,7 +1,7 @@
<?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE plist PUBLIC "-//Apple//DTD PLIST 1.0//EN" "http://www.apple.com/DTDs/PropertyList-1.0.dtd">
<!--
CB-504 — launchd agent for bridged (macOS).
CB-504 / CB-594 — launchd agent for bridged (macOS).
This is the real supervision target today: the dogfooded daemon runs on macOS, where there is
no systemd. A systemd unit ships alongside (deploy/bridged.service) for the Linux gateways
@@ -9,14 +9,33 @@
Install:
cp deploy/dev.ltms.bridged.plist ~/Library/LaunchAgents/
# edit the paths + JAVA_HOME below to match this host, then:
launchctl load -w ~/Library/LaunchAgents/dev.ltms.bridged.plist
launchctl list | grep bridged
The paths below are already filled in for this host (resolved 2026-08-16 from
`/usr/libexec/java_home`... except that reported the system Applet-plugin JVM, not the jenv-
managed JDK 25 actually used to build/run bridged, so JAVA_HOME here is the real one:
`JENV_VERSION=25.0.3 java -XshowSettings:properties -version 2>&1 | grep java.home`; `which mvn`;
`echo $HOME`). If this file is copied to a different host, re-resolve all three paths and check
no placeholder path is left behind; scripts/redeploy-bridged.sh's check mode does not (and
cannot) check this file for you.
CB-594 — launchd cannot run a login shell (see the PATH comment on EnvironmentVariables below,
and scripts/bridged-launchd-wrapper.sh for the fix): ProgramArguments below execs THAT wrapper,
not java directly, so WORKER_GITEA_TOKEN and AI_GATEWAY_TOKEN still get sourced from
${SHARED_ENV}/tools/secrets.sh even though launchd itself never sources anything.
Note on ordering: launchd has no "start after herdr" primitive for user agents, and neither
does systemd in a way that survives a socket appearing late. bridged retries the herdr socket
on startup instead, so an agent that comes up before herdr converges rather than dying — that
retry is the actual fix; KeepAlive below is the backstop.
CB-594 — KeepAlive vs. scripts/redeploy-bridged.sh: a bare SIGTERM makes this JVM exit 143 even
with its shutdown hook running to completion (measured, see the CB-594 report), which
SuccessfulExit:false below reads as a crash and races to restart the OLD jar. The redeploy
script now detects a loaded agent and uses `launchctl unload`/`load` instead of a raw kill, so
only one supervisor ever touches the process at a time — read that script's own output on a
redeploy for the confirmation.
-->
<plist version="1.0">
<dict>
@@ -25,22 +44,23 @@
<key>ProgramArguments</key>
<array>
<string>/Users/CHANGEME/Tool/jdk-25.0.2.jdk/Contents/Home/bin/java</string>
<string>/Users/dai.ha/LTMS/claude-bridge/scripts/bridged-launchd-wrapper.sh</string>
<string>/Users/dai.ha/Softwares/jdks/jdk-25.0.3.jdk/Contents/Home/bin/java</string>
<string>-jar</string>
<string>/Users/CHANGEME/src/claude-bridge/bridged/target/bridged.jar</string>
<string>/Users/dai.ha/LTMS/claude-bridge/bridged/target/bridged.jar</string>
<string>bridged.yaml</string>
</array>
<!-- Config path in ProgramArguments is relative, so the working directory must be the module. -->
<key>WorkingDirectory</key>
<string>/Users/CHANGEME/src/claude-bridge/bridged</string>
<string>/Users/dai.ha/LTMS/claude-bridge/bridged</string>
<key>EnvironmentVariables</key>
<dict>
<key>JAVA_HOME</key>
<string>/Users/CHANGEME/Tool/jdk-25.0.2.jdk/Contents/Home</string>
<string>/Users/dai.ha/Softwares/jdks/jdk-25.0.3.jdk/Contents/Home</string>
<key>HERDR_SOCKET_PATH</key>
<string>/Users/CHANGEME/.config/herdr/herdr.sock</string>
<string>/Users/dai.ha/.config/herdr/herdr.sock</string>
<!--
PATH matters more than it looks (CB-511): bridged propagates its own PATH to every worker
it spawns, so this line decides whether the fleet can run a build at all. launchd does NOT
@@ -48,11 +68,14 @@
bare /usr/bin:/bin and no JDK or Maven. Keep the toolchain entries first.
-->
<key>PATH</key>
<string>/Users/CHANGEME/Tool/jdk-25.0.2.jdk/Contents/Home/bin:/Users/CHANGEME/Tool/apache-maven-3.9.16/bin:/opt/homebrew/bin:/usr/local/bin:/usr/bin:/bin:/usr/sbin:/sbin</string>
<string>/Users/dai.ha/Softwares/jdks/jdk-25.0.3.jdk/Contents/Home/bin:/Users/dai.ha/Softwares/apache-maven/bin:/opt/homebrew/bin:/usr/local/bin:/usr/bin:/bin:/usr/sbin:/sbin</string>
<!--
Worker/API tokens are NOT set here: this file is committed. Export them from a private
launchd override or a wrapper script. bridged reads the API token from the env var named
by auth.tokenEnv (default BRIDGED_API_TOKEN) and only in auth.mode: token.
Worker/API tokens are NOT set here: this file is committed. CB-594 —
scripts/bridged-launchd-wrapper.sh (named in ProgramArguments above) is what supplies
them, by execing a login shell that sources ${SHARED_ENV}/tools/secrets.sh before the
daemon itself starts. bridged also reads the API token from the env var named by
auth.tokenEnv (default BRIDGED_API_TOKEN) and only in auth.mode: token — the wrapper
covers that one too, since it is the same login shell.
-->
</dict>
@@ -69,10 +92,17 @@
<key>ThrottleInterval</key>
<integer>10</integer>
<!--
CB-594 — same file scripts/redeploy-bridged.sh already tails ($BRIDGED/bridged.out), and both
streams point at it, not two separate log files: the script's fresh-line / ERROR-count checks
after a restart read this one path regardless of whether launchd or the script started the
process, and a stdout/stderr split would make half of what happens during a launchd-driven
restart invisible to it.
-->
<key>StandardOutPath</key>
<string>/Users/CHANGEME/src/claude-bridge/bridged/logs/bridged.out.log</string>
<string>/Users/dai.ha/LTMS/claude-bridge/bridged/bridged.out</string>
<key>StandardErrorPath</key>
<string>/Users/CHANGEME/src/claude-bridge/bridged/logs/bridged.err.log</string>
<string>/Users/dai.ha/LTMS/claude-bridge/bridged/bridged.out</string>
<key>ProcessType</key>
<string>Background</string>
+32
View File
@@ -0,0 +1,32 @@
#!/usr/bin/env bash
#
# CB-594 — the only reason this file exists: launchd does not run a login shell.
#
# WORKER_GITEA_TOKEN and AI_GATEWAY_TOKEN live in ${SHARED_ENV}/tools/secrets.sh, sourced only by a
# LOGIN shell (.zprofile/.zshrc etc). launchd execs a job's ProgramArguments directly — no shell, no
# profile, nothing sourced (the plist's own PATH comment documents the same gap one variable over).
# A daemon started that way boots fine and looks healthy; the failure is invisible until a worker
# tries to open a PR (WORKER_GITEA_TOKEN empty) or a gateway profile gets a 401 (AI_GATEWAY_TOKEN
# empty) — hours later, with nothing tying the two together (CB-591, CLAUDE.md "Redeploying the
# daemon"). Bridged now also logs which required secret names resolved at startup (see
# Bridged.reportRequiredSecrets), but that log line can only tell the truth if the tokens had a
# chance to be sourced in the first place — which is this script's entire job.
#
# So: launchd execs THIS script instead of java directly. This script execs a login shell
# ('zsh -l'), which sources secrets.sh, and that shell execs the real command in its place — one
# process throughout (exec, not a subshell fork), so launchd's PID tracking, KeepAlive, and
# StandardOut/ErrorPath all still see the one process they expect.
#
# The plist passes the full command as THIS script's own arguments, e.g.:
# ProgramArguments = [ .../bridged-launchd-wrapper.sh, /path/to/java, -jar, /path/to/bridged.jar,
# bridged.yaml ]
# so the wrapper stays generic and the actual command lives in exactly one place (the plist), not
# duplicated here.
set -euo pipefail
if [ "$#" -eq 0 ]; then
echo "bridged-launchd-wrapper.sh: no command given — check the plist's ProgramArguments" >&2
exit 2
fi
exec /bin/zsh -lc 'exec "$@"' -- "$@"
+72 -9
View File
@@ -5,13 +5,14 @@
# A merge is not a deployment: the running daemon holds the jar it was started with, so code merged
# to main does nothing until this runs. See CLAUDE.md -> "Redeploying the daemon".
#
# This script exists to turn five remembered traps into one auditable command:
# This script exists to turn six remembered traps into one auditable command:
#
# 1. A piped `mvn` hides BUILD FAILURE behind a zero exit, so the build here is never piped.
# 2. The daemon must start from a LOGIN shell, or the tokens it hands to members are empty:
# WORKER_GITEA_TOKEN (workers cannot open a PR) and AI_GATEWAY_TOKEN (401 at llm.ltms.dev).
# Both are read from the DAEMON's own environment at spawn time, so a value added to
# secrets.sh after startup is absent. Nothing logs this, so the script checks and says so.
# secrets.sh after startup is absent. Nothing logs this here, so the script checks and says
# so — and since CB-594, bridged's own startup log says so too, by env var name.
# 3. An old daemon that never actually died looks identical from the outside, so the script waits
# for the process to exit and for the port to free before it starts a new one.
# 4. "It started" is not "it works": the script polls /healthz until it answers, and reports the
@@ -19,6 +20,13 @@
# mismatch.
# 5. Restarting under live members drops their tickets, so the script refuses unless you confirm
# the fleet is drained.
# 6. CB-594 — the launchd agent (deploy/dev.ltms.bridged.plist), if installed and loaded, is a
# SECOND supervisor: its KeepAlive.SuccessfulExit=false restarts the daemon on any nonzero
# exit, and a bare SIGTERM makes this JVM exit 143 even with its shutdown hook running to
# completion (measured — see the CB-594 report). A plain `kill` here would race launchd's own
# restart of the OLD jar. So this script detects whether the agent is loaded and, only then,
# swaps `kill` + manual `nohup` for `launchctl unload`/`load` — the one supervisor in control
# at any moment is whichever one you asked to act, never both.
#
# Usage:
# scripts/redeploy-bridged.sh # build, confirm, restart, verify
@@ -43,13 +51,17 @@ HEALTH='http://127.0.0.1:8765/healthz'
STOP_WAIT=30 # seconds to wait for a clean exit before reporting failure
HEALTH_WAIT=60 # seconds to wait for /healthz to answer after start
# CB-594: the launchd agent this script must not fight with (see trap 6 above).
LAUNCHD_LABEL='dev.ltms.bridged'
LAUNCHD_PLIST="$HOME/Library/LaunchAgents/$LAUNCHD_LABEL.plist"
DO_BUILD=1; ASSUME_YES=0; CHECK_ONLY=0
for arg in "$@"; do
case "$arg" in
--yes|-y) ASSUME_YES=1 ;;
--no-build) DO_BUILD=0 ;;
--check) CHECK_ONLY=1 ;;
-h|--help) sed -n '3,30p' "${BASH_SOURCE[0]}"; exit 0 ;;
-h|--help) sed -n '3,37p' "${BASH_SOURCE[0]}"; exit 0 ;;
*) echo "unknown option: $arg (try --help)" >&2; exit 2 ;;
esac
done
@@ -61,6 +73,11 @@ die() { printf '\n FAIL %s\n\n' "$*" >&2; exit 1; }
jar_id() { [ -f "$JAR" ] && shasum -a 256 "$JAR" | cut -c1-12 || echo "absent"; }
running_pid() { pgrep -f "$PATTERN" || true; }
# `launchctl list <label>` exits 0 iff the label is loaded (registered with launchd) — true whether
# or not it is currently running, which is exactly "supervision is active" for our purposes. Read-
# only: neither helper below changes anything, so both are also safe under --check.
launchd_installed() { [ -f "$LAUNCHD_PLIST" ]; }
launchd_loaded() { launchctl list "$LAUNCHD_LABEL" >/dev/null 2>&1; }
# ---------------------------------------------------------------- report state
@@ -74,6 +91,22 @@ fi
ok "jar on disk: $(jar_id) ($([ -f "$JAR" ] && date -r "$JAR" '+%Y-%m-%d %H:%M:%S' || echo 'none'))"
ok "HEAD: $(git -C "$REPO" log --oneline -1)"
# CB-594: supervision state. Installed and loaded are different facts — a copied-but-never-loaded
# plist supervises nothing, and a loaded label with no file backing it (rare, but possible after an
# edited/moved plist) is still what launchd will act on.
if launchd_installed; then
ok "launchd agent installed: $LAUNCHD_PLIST"
else
warn "launchd agent NOT installed (no supervision — a crash will not restart the daemon)."
fi
SUPERVISED=0
if launchd_loaded; then
SUPERVISED=1
ok "launchd agent loaded ($LAUNCHD_LABEL) — launchd supervises this daemon"
else
warn "launchd agent not loaded — this script is the only thing that will restart the daemon."
fi
# The trap with no log line. Checked in a LOGIN shell, because that is how the daemon is started
# below. Never prints the value — only whether it resolved.
if zsh -lc '[ -n "${WORKER_GITEA_TOKEN:-}" ]' 2>/dev/null; then
@@ -137,11 +170,26 @@ if [ -n "$OLD_PID" ] && [ "$ASSUME_YES" = 0 ]; then
fi
# ------------------------------------------------------------------ stop
#
# CB-594: when SUPERVISED, launchd owns the stop — never a raw `kill` here. A bare SIGTERM makes
# this JVM exit 143 even with its shutdown hook running to completion (verified separately: a
# throwaway Java process with an equivalent shutdown hook, sent SIGTERM from a login shell that
# could `wait` on it directly, reported exit code 143 every time — never 0). launchd's
# KeepAlive.SuccessfulExit=false treats any nonzero exit as a crash and restarts the OLD jar,
# which would race this script's own restart of the NEW one. `launchctl unload` avoids that race
# by deregistering the job first, so no KeepAlive is left armed when the process actually stops.
if [ -n "$OLD_PID" ]; then
say "stop"
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)" # verify a FRESH line appears later
kill "$OLD_PID"
if [ "$SUPERVISED" = 1 ]; then
echo " supervision is ON: using 'launchctl unload' (not kill) so launchd's own KeepAlive"
echo " cannot restart the OLD jar out from under this script — see the CB-594 comment above."
launchctl unload -w "$LAUNCHD_PLIST" \
|| die "launchctl unload failed — the daemon may still be under supervision; investigate before retrying"
else
kill "$OLD_PID"
fi
for _ in $(seq "$STOP_WAIT"); do
[ -z "$(running_pid)" ] && break
sleep 1
@@ -152,18 +200,33 @@ if [ -n "$OLD_PID" ]; then
leave worktrees and panes behind. Investigate, then kill -9 by hand if you accept that."
fi
ok "pid $OLD_PID exited"
elif [ "$SUPERVISED" = 1 ]; then
# Loaded but not currently running (e.g. throttled after a crash loop). Unload it anyway so the
# start step below does a clean load, never a load stacked on an already-loaded label.
say "stop"
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)"
launchctl unload -w "$LAUNCHD_PLIST" 2>/dev/null || true
ok "launchd agent unloaded (was already not running)"
else
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)"
fi
# ------------------------------------------------------------------ start
# Login shell (zsh -l) is what puts the secrets on the daemon's environment. cwd must be bridged/
# because the daemon resolves bridged.yaml, logs/ and target/ relative to it.
# Unsupervised: login shell (zsh -l) is what puts the secrets on the daemon's environment, and cwd
# must be bridged/ because the daemon resolves bridged.yaml, logs/ and target/ relative to it.
# Supervised: launchd does both — deploy/dev.ltms.bridged.plist points ProgramArguments at
# scripts/bridged-launchd-wrapper.sh (CB-594), which is what execs the login shell in launchd's
# place, and WorkingDirectory in the plist already pins bridged/.
say "start"
# Absolute jar path so `ps` names which checkout is running; cwd still bridged/ because the daemon
# resolves bridged.yaml, logs/ and target/ relative to it.
( cd "$BRIDGED" && zsh -lc "nohup java -jar '$JAR' >> bridged.out 2>&1 &" )
if [ "$SUPERVISED" = 1 ]; then
echo " supervision is ON: using 'launchctl load' so launchd starts and keeps supervising this"
echo " process, instead of a manual nohup that launchd would know nothing about."
launchctl load -w "$LAUNCHD_PLIST" || die "launchctl load failed"
else
# Absolute jar path so `ps` names which checkout is running.
( cd "$BRIDGED" && zsh -lc "nohup java -jar '$JAR' >> bridged.out 2>&1 &" )
fi
for _ in $(seq 10); do
NEW_PID="$(running_pid)"