Compare commits

...

8 Commits

Author SHA1 Message Date
Dai Ha 16de9df000 CB-601: make the recovery-race test's head start deterministic, not a sleep
CI / contract (pull_request) Successful in 1m8s
CI / build (pull_request) Successful in 1m9s
2026-08-16 18:02:57 +02:00
ltms 613ece92dc Merge CB-594: make supervision and a working fleet possible at the same time
CI / contract (push) Successful in 49s
CI / build (push) Failing after 1m45s
The launchd unit was a CHANGEME template that had never been installed, and it could not have worked if it were: launchd does not source a login shell, so the daemon would have started with no forge or gateway token, and the failure would only appear much later as workers unable to open a PR.

Four parts: a wrapper that execs one login shell in place so the job inherits the secret store; a startup report naming which required secret env vars resolved and which are MISSING, by name only, never a value; the plist filled in with this host's real verified paths; and `redeploy-bridged.sh` detecting the agent and switching stop/start to `launchctl unload -w` / `load -w`, falling back to the existing kill + nohup when it is not installed.

Verified by the lead. My own unpiped build of the branch: 813 tests, BUILD SUCCESS. Trial-merged onto current main (which had moved twice) and rebuilt: Tests run: 822, Failures: 0, Errors: 0, Skipped: 0 — BUILD SUCCESS. All three host paths in the plist exist; no CHANGEME remains; the secret report leaks no values.

An independent reviewer, briefed only from the diff, verified the parts that matter today and said merge. It confirmed live that the unsupervised path is unchanged, that `launchctl list` exits 113 for an absent agent so the detection reads correctly, that the wrapper round-trips arguments containing spaces and quotes and stays a single exec chain, and that neither branch of the secret report can print a value.

Its two open findings only bite once the agent is actually loaded, which has not happened and is the operator's call. Filed as #91 — that must land before the agent is ever installed. The important one is that the script computes its log path from its own location while the plist hard-codes an absolute one; if they ever disagree, the post-restart ERROR check reads the wrong file and reports "ok" while the daemon crash-loops.

The agent is deliberately NOT installed by this merge. Nothing here changes how the daemon runs today.
2026-08-16 17:32:57 +02:00
ltms 8d4206c2b5 Merge CB-590: one nudge schedule per lead, with a reminder budget per source
CI / contract (push) Successful in 1m7s
CI / build (push) Failing after 1m15s
Carries PR #84 (its head `78ca24d` is an ancestor of this one), so this single merge delivers both rounds.

CB-590 collapses the CB-307 reply schedule and the CB-588 ticket schedule into one per lead, which closes the double-injection race. The review round found that this also collapsed the two reminder caps into one shared budget — so a busy reply stream could exhaust the cap and a ticket arriving afterwards would never be nudged at all. That was a regression, not a pre-existing wart: before this PR the two sources had independent counters.

The follow-up keeps the single schedule and gives each source its own budget. `decide()` returns INJECT while either source has pending work under its own cap, and STOP only when neither does. `stopOrRestart` and its snapshot-diff race logic are untouched.

Verified by the lead: `mvn -f bridged/pom.xml clean install` unpiped, exit code captured — Tests run: 814, Failures: 0, Errors: 0, Skipped: 0 — BUILD SUCCESS. `ReplyPushLoopTest`: 41 tests.

Known and deliberately out of scope: an item arriving during the 15s backoff is already inside the tick's "before" snapshot, so at-cap work can be abandoned rather than nudged once. That shape predates this PR in the ticket-only loop and is filed as #87 (CB-598).
2026-08-16 17:28:57 +02:00
ltms 03286a589b Merge CB-527/CB-528: bound AMQP prefetch, confirm publishes, and close the recovery race
CI / contract (push) Successful in 1m4s
CI / build (push) Failing after 1m18s
Carries PR #83 (its head is an ancestor of this one), so this single merge delivers both rounds.

CB-527 bounds prefetch. CB-528 makes a publish wait for its broker confirm, so a failed publish is never reported as success, and correlates a `mandatory` Return back to the right publish.

The review round on #83 found one real race: `failPendingPublishesOnRecovery` swept the pending maps without holding `publishChannelLock`, so a publish issued on the already-recovered channel could be failed by the sweep — the exact inversion of what CB-528 exists to prevent, and undetectable by dedup because `MessageService.reply` mints a fresh msgId per call. Fixed by taking the same lock. `close()` now fails in-flight publishes promptly instead of letting them time out after 10s.

Verified by the lead, not taken on the worker's word:
- `mvn -f bridged/pom.xml clean install` unpiped: Tests run: 809, Failures: 0, Errors: 0, Skipped: 0 — BUILD SUCCESS
- `mvn test -Pcontract -Dtest=AmqpReplyInboxContractTest` against a real broker: Tests run: 8 — BUILD SUCCESS

That second run mattered: #83's own new tests are contract-tagged, so the default suite and CI prove nothing about them.
2026-08-16 17:27:21 +02:00
Dai Ha 88b9503c3b CB-590 follow-up: give each nudge source its own reminder budget
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 1m46s
decide(lead, reminderCount) shared one counter across the reply and
ticket sources after PR #84 collapsed both onto a single per-lead
schedule. A reply stream that used up the whole budget could then
make decide() STOP even for a ticket that had never been nudged and
had coalesced onto the same still-active schedule — stranding it with
no live schedule left, since stopOrRestart's racedIn check does not
save work that was already present in the "before" snapshot.

decide() now tracks a per-source count (replyReminderCount,
ticketReminderCount) and returns INJECT while either source is still
under its own cap, STOP only when both are exhausted. Still exactly
one schedule per lead; stopOrRestart's snapshot-diff logic is
untouched.
2026-08-16 17:26:07 +02:00
Dai Ha c553d795d8 CB-528: close the recovery race in AmqpReplyInbox
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Failing after 1m20s
failPendingPublishesOnRecovery walked and cleared pendingBySeq/pendingByMsgId
without holding publishChannelLock, so a publish() that registered while the
sweep was still iterating could be failed even though it published on the
already-recovered channel — a successful publish reported as failed, and
since MessageService.reply mints a fresh msgId per retry, dedup can't catch
the resulting duplicate. Guard the sweep with publishChannelLock: publish()
only holds it for the seq/map-put/basicPublish, so the sweep can only ever
wait for an in-flight basicPublish to return, never a broker round trip.

Also make close() fail in-flight publishes immediately with a clear message
instead of leaving them to idle out the 10s confirm timeout, and record the
(currently unreachable) msgId-uniqueness assumption pendingByMsgId relies on.

AmqpReplyInboxRecoveryRaceTest drives the sweep and a real publish() against
each other directly (no live broker reconnect) using Proxy-backed fake AMQP
channels and a large in-flight backlog to make the race window observable;
confirmed it fails without the guard (reply "fresh" wrongly failed as
"connection recovered mid-publish") and passes with it.
2026-08-16 17:24:43 +02:00
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
Dai Ha 4fa6553db5 CB-527/CB-528: bound AMQP prefetch and confirm publishes before claiming durable
CI / build (pull_request) Successful in 1m3s
CI / contract (pull_request) Successful in 1m15s
CB-527: basicQos(prefetch) on the consume channel before basicConsume, configurable
via broker.prefetch (default 32), so an undrained inbox backlog stays on the broker
instead of growing the JVM heap without limit.

CB-528: publish moves to its own confirm-mode channel with mandatory=true and a
return listener, so an unroutable or unconfirmed reply now throws instead of
vanishing silently. The confirm callback checks the per-message returned flag
(set by the return listener, which the broker always fires before the matching
confirm) so an acked-but-returned publish is still reported as a failure. The ack
path stays on its own channel/lock and never waits on a publish confirm.
2026-08-16 16:55:40 +02:00
12 changed files with 1092 additions and 81 deletions
+6 -2
View File
@@ -434,10 +434,14 @@ 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.
# 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.
# 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,8 +352,9 @@ 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());
log.info("reply inbox: AMQP broker (durable) at {}", cfg.broker().uri());
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)");
@@ -543,6 +548,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,6 +5,7 @@ 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;
@@ -532,16 +533,24 @@ 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 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}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Broker(String uri) {
public record Broker(String uri, Integer prefetch) {
/** 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,6 +7,7 @@ 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;
@@ -14,7 +15,13 @@ 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
@@ -29,6 +36,27 @@ import java.util.concurrent.ConcurrentHashMap;
* 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
@@ -54,6 +82,12 @@ 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. */
@@ -63,39 +97,87 @@ 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. Assumes {@code msgId} is unique per in-flight publish: a second {@link #publish} for a
* {@code msgId} still awaiting its confirm would overwrite this entry and misdirect
* {@link #onReturn}'s lookup. Not reachable today — {@code MessageService.reply} generates a fresh
* {@code UUID} per call — so no guard is added for it.
*/
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) {}
/** Connect to {@code uri} (e.g. {@code amqp://guest:guest@127.0.0.1:5672/}) and open the inbox. */
/** 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}. */
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"));
return new AmqpReplyInbox(factory.newConnection("bridged-reply-inbox"), prefetch);
} catch (Exception e) {
throw new IllegalStateException("cannot connect to AMQP broker at " + uri, e);
}
}
/** Wrap an already-open connection (injection seam for the contract test). */
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (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).
// 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.
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");
}
@@ -141,6 +223,12 @@ 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()
@@ -148,12 +236,35 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
.deliveryMode(2) // persistent — survives a broker restart
.contentType("text/plain")
.build();
try {
synchronized (channelLock) {
channel.basicPublish("", queueName(target), props, content.getBytes(StandardCharsets.UTF_8));
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);
}
} catch (IOException e) {
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);
}
}
@@ -222,17 +333,115 @@ 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.
* Guarded by {@link #publishChannelLock}, the same lock {@link #publish} holds while it takes its
* sequence number and registers its {@link Pending}: without it, a {@link #publish} that starts
* after the connection has already recovered (so it publishes — and will be confirmed — on the
* <em>new</em> channel) can register between this sweep's iteration and its clear, and this sweep
* then fails a publish that actually succeeded. {@link #publish} only holds the lock for the
* seq/map-put/{@code basicPublish} — it awaits the confirm outside it — so this sweep can only ever
* wait for an in-flight {@code basicPublish} call to return, never for a broker round trip. No
* deadlock.
*
* <p>Package-private (rather than {@code private}) only so the unit test can drive it directly
* against a concurrent {@link #publish} without a live broker reconnect.
*/
void failPendingPublishesOnRecovery() {
synchronized (publishChannelLock) {
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
pending.confirmed.completeExceptionally(new IllegalStateException(
"AMQP connection recovered mid-publish; confirm status of reply " + pending.msgId
+ " is unknown"));
}
}
}
/**
* Fail every publish still awaiting its confirm with a clear, immediate error instead of leaving it
* to time out after {@link #CONFIRM_TIMEOUT_MS} once the channels are closed underneath it. Guarded
* by {@link #publishChannelLock} for the same reason as {@link #failPendingPublishesOnRecovery}.
*/
private void failPendingPublishesOnClose() {
synchronized (publishChannelLock) {
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
pending.confirmed.completeExceptionally(new IllegalStateException(
"AMQP reply inbox closed while publish of reply " + pending.msgId
+ " was still awaiting its confirm"));
}
}
}
private static String queueName(String target) {
return QUEUE_PREFIX + target + QUEUE_SUFFIX;
}
@Override
public void close() {
failPendingPublishesOnClose();
try {
channel.close();
} 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) {
@@ -37,10 +37,12 @@ import java.util.stream.Collectors;
* <p>Each tick examines <em>everything</em> pending for that lead — reply targets whose inbox
* still holds an unacked message ({@link #pendingReplies}) and tickets not yet collected
* ({@link #pendingTickets}) — and sends at most one combined nudge per tick
* ({@link #injectNudge(String, int)}). Work that arrives while the lead is busy is never lost:
* it is re-read fresh on every tick until the lead is injectable or the shared reminder cap
* ({@link #maxReminders}) is reached, whichever the durable inbox / pending-ticket set doesn't
* already answer via {@code STOP}.
* ({@link #injectNudge(String, int, int)}). Work that arrives while the lead is busy is never
* lost: it is re-read fresh on every tick until the lead is injectable or its own reminder cap
* ({@link #maxReminders}) is reached — reply and ticket work each spend from their own budget, so
* one source exhausting its cap does not stop nudges about the other (post-CB-590 regression fix;
* see {@link #decide}) — whichever the durable inbox / pending-ticket set doesn't already answer
* via {@code STOP}.
*/
public final class ReplyPushLoop {
@@ -144,20 +146,32 @@ public final class ReplyPushLoop {
* Pure decision function: examine everything pending for {@code lead} — reply targets and
* tickets alike — and return what the loop should do.
*
* @param lead the lead terminal to nudge
* @param reminderCount how many nudges have been sent so far for this lead (shared across
* both reply and ticket work — CB-590 collapses the reminder cap onto
* one counter per lead, so alternating sources cannot outrun the bound)
* <p><strong>CB-590-fix: one schedule, two budgets.</strong> The single per-lead schedule
* (CB-590) still ticks once for both sources, but each source is capped independently —
* {@code replyReminderCount} against a reply target still pending, {@code ticketReminderCount}
* against a ticket still pending. A busy reply stream that exhausts its own cap must not stop
* the loop from nudging about a ticket that still has budget left, and vice versa: either
* source being eligible (has pending work AND is under its own cap) is enough for
* {@link Action#INJECT}. Only when neither source has eligible work does the loop
* {@link Action#STOP}.
*
* @param lead the lead terminal to nudge
* @param replyReminderCount how many nudges have covered pending reply work for this lead
* @param ticketReminderCount how many nudges have covered pending ticket work for this lead
* @return the action the caller should take
*/
Action decide(String lead, int reminderCount) {
boolean anyPending = !pendingReplyTargetsFor(lead).isEmpty() || !pendingTicketIdsFor(lead).isEmpty();
if (!anyPending) {
Action decide(String lead, int replyReminderCount, int ticketReminderCount) {
boolean hasReplyWork = !pendingReplyTargetsFor(lead).isEmpty();
boolean hasTicketWork = !pendingTicketIdsFor(lead).isEmpty();
if (!hasReplyWork && !hasTicketWork) {
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
return Action.STOP;
}
if (reminderCount >= maxReminders) {
log.debug("push: reminder cap ({}) reached for lead {}, stopping", maxReminders, lead);
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
boolean ticketEligible = hasTicketWork && ticketReminderCount < maxReminders;
if (!replyEligible && !ticketEligible) {
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
maxReminders, lead);
countNudge("exhausted");
return Action.STOP;
}
@@ -241,21 +255,28 @@ public final class ReplyPushLoop {
return;
}
log.debug("push: starting reminder loop for lead {}", lead);
scheduleNext(lead, 0);
scheduleNext(lead, 0, 0);
}
/** Execute one loop tick — called on the scheduler thread. */
private void tick(String lead, int reminderCount) {
private void tick(String lead, int replyReminderCount, int ticketReminderCount) {
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
Set<String> ticketsBefore = pendingTicketIdsFor(lead);
var action = decide(lead, reminderCount);
var action = decide(lead, replyReminderCount, ticketReminderCount);
switch (action) {
case INJECT -> {
injectNudge(lead, reminderCount);
scheduleNext(lead, reminderCount + 1);
injectNudge(lead, replyReminderCount, ticketReminderCount);
// Only the source(s) actually eligible this tick spend a unit of their own budget —
// an exhausted source riding along in the combined message (still pending, still
// named) does not get charged again; its count stays put until it drains.
boolean replyEligible = !repliesBefore.isEmpty() && replyReminderCount < maxReminders;
boolean ticketEligible = !ticketsBefore.isEmpty() && ticketReminderCount < maxReminders;
scheduleNext(lead,
replyEligible ? replyReminderCount + 1 : replyReminderCount,
ticketEligible ? ticketReminderCount + 1 : ticketReminderCount);
}
// Re-check after the configured backoff; the lead may become injectable soon.
case WAIT_BUSY -> scheduleNext(lead, reminderCount);
case WAIT_BUSY -> scheduleNext(lead, replyReminderCount, ticketReminderCount);
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore);
}
}
@@ -296,14 +317,14 @@ public final class ReplyPushLoop {
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t));
if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
scheduleNext(lead, 0);
scheduleNext(lead, 0, 0);
return;
}
log.debug("push: reminder loop ended for lead {}", lead);
}
/** Send one combined nudge covering everything currently pending for {@code lead}. */
private void injectNudge(String lead, int reminderCount) {
private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount) {
// Re-read rather than threading it down from decide(): a reply can drain, or a ticket be
// collected (or another arrive), between the decision and the injection.
Set<String> replyTargets = pendingReplyTargetsFor(lead);
@@ -315,18 +336,20 @@ public final class ReplyPushLoop {
String nudge = formatNudge(replyTargets, tickets);
try {
agents.send(lead, nudge);
log.debug("push: nudge {}/{} sent to lead {} ({} reply target(s), {} ticket(s))",
reminderCount + 1, maxReminders, lead, replyTargets.size(), tickets.size());
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}; {} reply target(s), {} ticket(s))",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
replyTargets.size(), tickets.size());
countNudge("delivered");
} catch (RuntimeException e) {
log.warn("push: failed to nudge lead {} (reminder {}/{}): {}",
lead, reminderCount + 1, maxReminders, e.toString());
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}): {}",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, e.toString());
}
}
/** Schedule the next tick on the scheduler thread pool. */
private void scheduleNext(String lead, int nextReminderCount) {
scheduler.schedule(() -> tick(lead, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS);
private void scheduleNext(String lead, int nextReplyReminderCount, int nextTicketReminderCount) {
scheduler.schedule(() -> tick(lead, nextReplyReminderCount, nextTicketReminderCount),
backoffMs, TimeUnit.MILLISECONDS);
}
// --- nudge formatting ------------------------------------------------------------------------
@@ -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,6 +16,7 @@ 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;
/**
@@ -166,6 +167,88 @@ 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)
@@ -0,0 +1,288 @@
package dev.ltms.bridged.msg;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.ConfirmCallback;
import com.rabbitmq.client.Connection;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
import java.lang.reflect.InvocationHandler;
import java.lang.reflect.Proxy;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* CB-528 follow-up: {@link AmqpReplyInbox#failPendingPublishesOnRecovery()} must not fail a publish
* that registers concurrently with (and was not yet visible when) the recovery sweep began — that
* would report a publish that actually succeeded as failed, and {@code MessageService.reply} retries
* with a fresh {@code msgId}, so the reply is delivered twice. A live broker reconnect cannot be
* forced reliably, so this drives {@link AmqpReplyInbox#failPendingPublishesOnRecovery} and a real
* {@link AmqpReplyInbox#publish} against each other directly, against fake AMQP channels built with
* {@link Proxy} (no mocking library is on the classpath).
*
* <p>Also covers Finding 2 (CB-528 follow-up): {@link AmqpReplyInbox#close()} must fail an in-flight
* publish promptly instead of leaving it to idle out the 10s confirm timeout.
*/
class AmqpReplyInboxRecoveryRaceTest {
/** Large enough that thousands of entries are still unprocessed by the time the very first one
* is observed as failed (see {@code sweepIsHoldingTheLock} below) — that gap is what makes the
* head start deterministic instead of a coin flip. 100,000 gave the same guarantee but made the
* test far more expensive than the guarantee needs; the ordering no longer depends on a timing
* window sized to the full backlog; just to the tail of it. */
private static final int STALE_PUBLISHES = 2_000;
@Test
@Timeout(30)
void recoverySweepDoesNotFailAPublishThatRegistersWhileItIsRunning() throws Exception {
AtomicLong seqCounter = new AtomicLong();
List<Long> seqOrder = new CopyOnWriteArrayList<>();
List<String> msgIdOrder = new CopyOnWriteArrayList<>();
AtomicReference<ConfirmCallback> ackCallback = new AtomicReference<>();
AtomicReference<ConfirmCallback> nackCallback = new AtomicReference<>();
Channel publishChannel = fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback);
Channel consumeChannel = fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback);
Connection connection = fakeConnection(consumeChannel, publishChannel);
AmqpReplyInbox inbox = new AmqpReplyInbox(connection, AmqpReplyInbox.DEFAULT_PREFETCH);
// STALE_PUBLISHES in-flight publishes that never get confirmed — they sit in pendingBySeq /
// pendingByMsgId exactly like publishes whose confirm never arrived before a connection drop.
// Virtual threads make this many concurrent blocking publish() calls cheap.
CountDownLatch staleStarted = new CountDownLatch(STALE_PUBLISHES);
// Counted down by the FIRST stale publish thread to observe its own failure. That can only
// happen from inside failPendingPublishesOnRecovery() — nothing else in this test ever
// completes a stale Pending exceptionally (no nack/return is simulated for any "stale-*"
// msgId) — so seeing it fire is direct, observable proof the sweep is inside its loop, not a
// timing guess. It replaces the old fixed Thread.sleep(5) head start.
CountDownLatch sweepIsHoldingTheLock = new CountDownLatch(1);
for (int i = 0; i < STALE_PUBLISHES; i++) {
String msgId = "stale-" + i;
Thread.ofVirtual().start(() -> {
staleStarted.countDown();
try {
inbox.publish("worker-stale", msgId, "x");
} catch (IllegalStateException expected) {
sweepIsHoldingTheLock.countDown();
}
});
}
staleStarted.await();
// Let the registrations (the synchronized put into pendingBySeq/pendingByMsgId) actually land
// for all of them before the sweep starts, so the sweep begins with a large, real backlog.
Thread.sleep(300);
AtomicReference<Throwable> sweepError = new AtomicReference<>();
Thread sweepThread = new Thread(() -> {
try {
inbox.failPendingPublishesOnRecovery();
} catch (Throwable t) {
sweepError.set(t);
}
}, "recovery-sweep");
sweepThread.start();
// Deterministic head start: block until the sweep has actually failed one of the stale
// publishes. failPendingPublishesOnRecovery() (once guarded, as it is on main) holds
// publishChannelLock for its ENTIRE loop, not just per entry — so this failure proves the
// sweep is, at this instant, still holding that lock. With STALE_PUBLISHES this large, the
// remaining ~1,999 entries give an enormous margin between "first failure observed" and "sweep
// releases the lock": there is no window left for "fresh" to slip in before the sweep starts,
// or to win the lock ahead of it — see the case-2 note below. This also means Case 1 (the sweep
// is already inside its loop, holding the lock, when "fresh" tries to register) is now
// guaranteed by construction rather than merely likely under a fixed sleep.
assertTrue(sweepIsHoldingTheLock.await(20, TimeUnit.SECONDS),
"the sweep never failed a single stale publish — it may not have started");
// Case 2 ("fresh" wins publishChannelLock before the sweep even starts, so it genuinely
// published on the stale channel and the sweep correctly fails it) is impossible by
// construction in this test: freshThread.start() below is reached only after
// sweepIsHoldingTheLock has counted down, which can only happen once
// failPendingPublishesOnRecovery() is already running and has already failed a stale entry.
// There is no code path that lets "fresh" start before the sweep starts. That case is real
// and correct production behaviour (see AmqpReplyInbox#failPendingPublishesOnRecovery's
// javadoc), it is just not reachable from this deterministic ordering, so it does not need a
// separate assertion here.
// This is the exact interleaving CB-528's follow-up describes: "the still-running recovery
// sweep" racing a publish that registers while it is mid-flight.
AtomicReference<Throwable> publishError = new AtomicReference<>();
Thread freshThread = new Thread(() -> {
try {
inbox.publish("worker-fresh", "fresh", "hello");
} catch (Throwable t) {
publishError.set(t);
}
}, "fresh-publish");
freshThread.start();
sweepThread.join(20_000);
assertNull(sweepError.get(), "sweep threw: " + sweepError.get());
// Simulate the broker's real confirm for "fresh" now that the sweep is done, so a correct
// implementation's publish() returns normally instead of idling out CONFIRM_TIMEOUT_MS. Poll
// for the registration rather than checking once: freshThread may still be contending for
// publishChannelLock (behind the sweep's own hold on it, and possibly other stale threads
// still unwinding) even though the sweep itself has already finished.
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(9);
int idx = -1;
while (idx < 0 && System.nanoTime() < deadline) {
idx = msgIdOrder.indexOf("fresh");
if (idx < 0) {
Thread.sleep(20);
}
}
if (idx >= 0 && ackCallback.get() != null) {
ackCallback.get().handle(seqOrder.get(idx), false);
}
freshThread.join(15_000);
assertNull(publishError.get(),
"a publish that registered while the recovery sweep was running must not be failed by "
+ "it, but got: " + publishError.get());
}
@Test
@Timeout(15)
void closeFailsInFlightPublishPromptlyInsteadOfWaitingOutTheConfirmTimeout() throws Exception {
AtomicLong seqCounter = new AtomicLong();
List<Long> seqOrder = new CopyOnWriteArrayList<>();
List<String> msgIdOrder = new CopyOnWriteArrayList<>();
AtomicReference<ConfirmCallback> ackCallback = new AtomicReference<>();
AtomicReference<ConfirmCallback> nackCallback = new AtomicReference<>();
Channel publishChannel = fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback);
Channel consumeChannel = fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback);
Connection connection = fakeConnection(consumeChannel, publishChannel);
AmqpReplyInbox inbox = new AmqpReplyInbox(connection, AmqpReplyInbox.DEFAULT_PREFETCH);
CountDownLatch publishReturned = new CountDownLatch(1);
AtomicReference<Throwable> publishError = new AtomicReference<>();
AtomicLong elapsedMillis = new AtomicLong();
Thread publishThread = new Thread(() -> {
long start = System.nanoTime();
try {
inbox.publish("worker-close", "never-confirmed", "x");
} catch (Throwable t) {
publishError.set(t);
} finally {
elapsedMillis.set((System.nanoTime() - start) / 1_000_000);
publishReturned.countDown();
}
}, "publish-during-close");
publishThread.start();
Thread.sleep(200); // let publish() register before close() runs
inbox.close();
assertTrue(publishReturned.await(5, TimeUnit.SECONDS), "publish() did not return after close()");
assertNotNull(publishError.get(), "a publish in flight when close() runs must fail, not hang");
assertTrue(publishError.get().getMessage() != null
&& publishError.get().getMessage().toLowerCase().contains("closed"),
"expected a clear closed-inbox message, got: " + publishError.get());
assertTrue(elapsedMillis.get() < 5_000,
"close() should fail the in-flight publish promptly, not wait out the confirm timeout — took "
+ elapsedMillis.get() + "ms");
}
/** A {@link Proxy}-backed {@link Channel}: only the calls {@link AmqpReplyInbox} actually makes
* are meaningfully implemented; everything else returns a harmless default. */
private static Channel fakeChannel(AtomicLong seqCounter, List<Long> seqOrder, List<String> msgIdOrder,
AtomicReference<ConfirmCallback> ackCallback,
AtomicReference<ConfirmCallback> nackCallback) {
InvocationHandler handler = (proxy, method, args) -> {
String name = method.getName();
if (name.equals("getNextPublishSeqNo")) {
long value = seqCounter.incrementAndGet();
seqOrder.add(value);
return value;
}
if (name.equals("basicPublish")) {
AMQP.BasicProperties props = (AMQP.BasicProperties) args[3];
msgIdOrder.add(props.getMessageId());
return null;
}
if (name.equals("addConfirmListener")) {
ackCallback.set((ConfirmCallback) args[0]);
nackCallback.set((ConfirmCallback) args[1]);
return null;
}
if (name.equals("equals")) {
return proxy == args[0];
}
if (name.equals("hashCode")) {
return System.identityHashCode(proxy);
}
if (name.equals("toString")) {
return "FakeChannel";
}
return defaultValue(method.getReturnType());
};
return (Channel) Proxy.newProxyInstance(AmqpReplyInboxRecoveryRaceTest.class.getClassLoader(),
new Class<?>[] {Channel.class}, handler);
}
/** A {@link Proxy}-backed {@link Connection} handing out {@code first} then {@code second} from
* successive {@code createChannel()} calls, matching {@link AmqpReplyInbox}'s constructor. */
private static Connection fakeConnection(Channel first, Channel second) {
AtomicInteger calls = new AtomicInteger();
InvocationHandler handler = (proxy, method, args) -> {
String name = method.getName();
if (name.equals("createChannel") && (args == null || args.length == 0)) {
return calls.getAndIncrement() == 0 ? first : second;
}
if (name.equals("equals")) {
return proxy == args[0];
}
if (name.equals("hashCode")) {
return System.identityHashCode(proxy);
}
if (name.equals("toString")) {
return "FakeConnection";
}
return defaultValue(method.getReturnType());
};
return (Connection) Proxy.newProxyInstance(AmqpReplyInboxRecoveryRaceTest.class.getClassLoader(),
new Class<?>[] {Connection.class}, handler);
}
private static Object defaultValue(Class<?> type) {
if (!type.isPrimitive() || type == void.class) {
return null;
}
if (type == boolean.class) {
return Boolean.FALSE;
}
if (type == long.class) {
return 0L;
}
if (type == short.class) {
return (short) 0;
}
if (type == byte.class) {
return (byte) 0;
}
if (type == char.class) {
return (char) 0;
}
if (type == double.class) {
return 0.0d;
}
if (type == float.class) {
return 0.0f;
}
return 0;
}
}
@@ -81,7 +81,7 @@ class ReplyPushLoopTest {
@Test
void decideWithNothingPendingIsStop() {
agents = agentWithStatus("idle");
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0));
}
@Test
@@ -90,7 +90,7 @@ class ReplyPushLoopTest {
inbox.publish(WORKER, "m1", "hello");
var loop = loop(2, 100);
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0));
}
@Test
@@ -99,7 +99,7 @@ class ReplyPushLoopTest {
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
}
@Test
@@ -108,7 +108,7 @@ class ReplyPushLoopTest {
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0),
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
"BLOCKED is injectable");
}
@@ -118,7 +118,7 @@ class ReplyPushLoopTest {
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0),
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
"DONE is injectable");
}
@@ -128,7 +128,7 @@ class ReplyPushLoopTest {
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
}
@Test
@@ -137,7 +137,7 @@ class ReplyPushLoopTest {
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
}
@Test
@@ -146,9 +146,9 @@ class ReplyPushLoopTest {
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
inbox.ack(WORKER, "m1");
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0));
}
// --- onReplyQueued integration -------------------------------------------------------------
@@ -240,7 +240,7 @@ class ReplyPushLoopTest {
@Test
void decideTicketsWithNothingPendingIsStop() {
agents = agentWithStatus("idle");
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0));
}
@Test
@@ -248,7 +248,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
var loop = loop(2, 100_000);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 2));
}
@Test
@@ -256,7 +256,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
var loop = loop(5, 100_000);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
}
@Test
@@ -264,7 +264,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("working");
var loop = loop(5, 100_000);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
}
@Test
@@ -272,7 +272,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 100_000);
loop.onTicketTerminal("task-1", WORKER, false); // no lead known -> never registered as pending
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0));
}
// --- CB-588: onTicketTerminal integration ---------------------------------------------------
@@ -420,6 +420,49 @@ class ReplyPushLoopTest {
+ "restart the loop — that would defeat the reminder cap");
}
// --- mirror of the two stopOrRestart tests above, for the reply arm ------------------------
//
// Both tests above only ever passed Set.of() for repliesBefore, so racedIn's reply branch
// (`pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t))`) was
// never exercised by anything other than an always-empty snapshot. The reviewer flagged this:
// racedIn is symmetric in the code, and only half of it was pinned by a test.
@Test
void aReplyStillPendingWhenTheLoopStopsIsNotStranded() {
// Mirrors aTicketStillPendingWhenTheLoopStopsIsNotStranded: a reply target that raced in
// during the decision-to-release window (absent from the "before" snapshot) must reclaim
// the schedule slot rather than being stranded with no schedule left to nudge about it.
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
ReplyPushLoop loop = loop(1, 100_000); // long backoff — no natural tick fires during this test
loop.onReplyQueued(WORKER); // pendingReplies={term_worker}; activeLeads={PRIMARY}
loop.stopOrRestart(PRIMARY, Set.of(), Set.of());
assertTrue(loop.isActive(), "a reply that raced the loop's stop must reclaim the schedule "
+ "slot, not be stranded with no schedule left to ever nudge about it");
}
@Test
void aStaleUncollectedReplyAtCapDoesNotRestartTheLoop() {
// Mirrors aStaleUncollectedTicketAtCapDoesNotRestartTheLoop: a reply target already
// accounted for at decide time (present in repliesBefore) must not restart the loop —
// that is the reminder cap doing its job, not a race.
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
ReplyPushLoop loop = loop(1, 100_000);
loop.onReplyQueued(WORKER); // pendingReplies={term_worker}; activeLeads={PRIMARY}
loop.stopOrRestart(PRIMARY, Set.of(WORKER), Set.of());
assertFalse(loop.isActive(), "a stale reply target already accounted for at decide time "
+ "must not restart the loop — that would defeat the reminder cap");
}
@Test
void ticketNudgeFormatIsCorrect() {
String single = ReplyPushLoop.TICKET_NUDGE_FORMAT.formatted("task-1", "", "task-1");
@@ -485,6 +528,57 @@ class ReplyPushLoopTest {
assertTrue(nudge.contains("task-1"), "the ticket must not be dropped: " + nudge);
}
// --- CB-590 follow-up: per-source reminder budgets — the regression this round exists for ---
@Test
void oneExhaustedSourceDoesNotBlockANudgeForTheOtherSource() {
// The live trace this ticket was filed from: an undrained reply target got nudged up to
// its cap (5 reminders), then a ticket for the SAME lead went terminal shortly before the
// next scheduled tick — so it coalesced onto the still-active schedule (arriving BEFORE
// that tick's "before" snapshot, not during the decision-to-release race stopOrRestart
// guards). With CB-590's single shared reminder counter, that tick's decide() saw
// reminderCount already at the cap and returned STOP regardless of the ticket, and because
// the ticket was already present in that tick's "before" snapshot, stopOrRestart's
// racedIn check (proven correct on its own above) did not save it either — it is not a
// race, it looks like ordinary stale backlog. The ticket was then stranded: pending
// forever with no live schedule, never named in any nudge.
//
// Fixed by giving each source its own counter. Here the reply source is AT its cap (2/2)
// and the ticket source has NEVER been nudged (0/2) — decide() must still return INJECT,
// because the ticket is still eligible on its own budget.
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
var loop = loop(2, 100_000); // huge backoff — this test drives decide()/isActive() directly
loop.onReplyQueued(WORKER);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 2, 0),
"the reply source is exhausted (2/2), but the ticket source has never been "
+ "nudged (0/2) — the lead must still be injected so the ticket is not "
+ "lost, exactly the CB-590 follow-up regression");
// isActive() (criterion #5): a real tick that takes the INJECT branch above never calls
// stopOrRestart, so the schedule started by onTicketTerminal above stays live — the
// ticket is not left stranded with isActive()==false while it is still pending.
assertTrue(loop.isActive(), "the schedule must stay active while the ticket source still "
+ "has budget left, even though the reply source sharing it is exhausted");
}
@Test
void bothSourcesExhaustedIsStillStop() {
// The flip side: per-source budgets must not turn into unbounded nudging. When BOTH
// sources are at their cap, decide() must still STOP — a per-source budget is still a
// budget.
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
var loop = loop(2, 100_000);
loop.onReplyQueued(WORKER);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 2),
"both the reply and the ticket source are at their own cap — must still stop");
}
// --- metrics (CB-512) ----------------------------------------------------------------------
@Test
@@ -513,7 +607,7 @@ class ReplyPushLoopTest {
var loop = loop(2, 100, metrics);
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0));
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the reminder cap must count as exhausted");
@@ -541,7 +635,7 @@ class ReplyPushLoopTest {
var loop = loop(2, 100_000, metrics);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 2));
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the ticket reminder cap must count as exhausted");
+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)"