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
12 changed files with 626 additions and 1155 deletions
+8 -41
View File
@@ -88,28 +88,15 @@ bind:
# backoffMs: 60000
# quietNudgeCap: 3
# Fleet health detection is dormant unless enabled (CB-573). It reads one whole-fleet agent list
# per tick.
# intervalSeconds → how often a tick runs (default 30). ENFORCED floor of 15: the code computes
# Math.max(15, intervalSeconds), so a lower value is silently raised, not
# rejected.
# workingSuspectAfterSeconds, paneProbeIntervalSeconds → accepted and parsed, but NOT YET READ by
# anything — the dormant monitor only consumes intervalSeconds today (CB-573
# shipped ahead of the evidence publishers these two knobs are for). Setting
# them changes nothing right now, and no minimum is enforced on either, because
# nothing reads them to enforce one. They exist so a later build can start
# honouring them without another config-shape change.
# notifications.mode → "webhook" flips what bridge_list REPORTS (healthCoverage: "full" instead
# of "detection-only") — it does NOT make bridged send any webhook call; no
# delivery mechanism is implemented yet. Any other value, or omitting the
# block, reports "detection-only".
# Fleet health detection is dormant unless enabled. It reads one whole-fleet agent list per tick.
# It can run without a webhook; bridge_list then reports healthCoverage: detection-only.
# health:
# enabled: true
# intervalSeconds: 30
# workingSuspectAfterSeconds: 600
# paneProbeIntervalSeconds: 60
# intervalSeconds: 30 # minimum 15
# workingSuspectAfterSeconds: 600 # minimum 300
# paneProbeIntervalSeconds: 60 # minimum 60
# notifications:
# mode: disabled
# mode: disabled # disabled (default) or webhook
# herdr Unix socket. Omit to use the client default
# (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}).
@@ -299,11 +286,6 @@ placement: weighted
# / credentialId. Those are hot because the placement policy (and, for credentialId,
# the CB-578 stage B quarantine check) reads them through a supplier — being config is
# not by itself enough to make a key hot.
# EXCEPT `fleet.leaders`: Bridged.main reads it once at startup to build the lead tab
# scanner and launcher, and neither is rebuilt on reload. A changed/added/removed
# `fleet.leaders` entry is silently accepted — the reload reports "config reloaded"
# with nothing in the deferred list — but has NO effect until you restart. Treat it
# as deferred in practice, even though today's reload output does not say so.
# DEFERRED → accepted into the new config, but the wiring built at startup keeps the old value
# until you restart: `lifecycle:`, `leadHeartbeat:`, `guard:`, `worktreeRoot:`,
# `spawnReadyTimeoutMs` / `spawnReadyPollMs`, `quarantineCooldownSeconds` (CB-578
@@ -386,12 +368,6 @@ fleet:
# An auto-launched lead is NOT a member: it gets no worker reply charter, is never registered with
# the session lifecycle (the idle reaper would kill your orchestrator), and stays on the
# subscription — ANTHROPIC_BASE_URL/AUTH_TOKEN are stripped from its env whatever the profile says.
#
# GET THE `tab:` VALUE RIGHT. A pane that does not match any configured `tab:` (a typo, a renamed
# tab, a pane no entry names at all) is not recognised as a lead — it resolves as an ordinary
# WORKER instead, silently, and every orchestration call it makes (spawn/stop/send/drain) is
# refused. There is no error at startup for this: an unmatched pane is simply not a lead. If your
# primary suddenly can't spawn or send, check this section first.
# leaders:
# opus-5.0:
# profile: opus # omit to never create this lead, only recognise it
@@ -447,15 +423,10 @@ guard:
# idleTtlSeconds → reap READY/DONE sessions idle longer than this (never BUSY/SPAWNING)
# contextCap → force-release a session after this many delegated turns
# drainTimeoutSeconds → seconds to wait for BUSY sessions on shutdown before forced teardown
# clearAfterTurn → whether a reusable worker discards its conversation context after every
# completed delegated turn (default false). Works for claude-code workers
# only — any other peer kind (e.g. opencode) logs "context reset is
# unsupported for peer kind …" once and the reset is a no-op.
# lifecycle:
# idleTtlSeconds: 300
# contextCap: 10
# drainTimeoutSeconds: 5
# clearAfterTurn: false
# Durable reply delivery (CB-307 Stage 2). OMIT this block entirely to keep the default
# in-memory, soft-state reply inbox (late worker replies are held only until a daemon bounce).
@@ -463,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,87 +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. 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) {}
/** 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");
}
@@ -223,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()
@@ -236,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);
}
}
@@ -333,115 +222,17 @@ 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) {
@@ -8,8 +8,6 @@ import dev.ltms.bridged.metrics.Metrics;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
@@ -18,43 +16,35 @@ import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
/**
* A status-gated push loop that nudges a lead's own herdr pane when it has uncollected work
* waiting: a worker reply queued with no live {@code bridge_send} to resolve it (CB-307), or an
* async delegation ticket ({@code bridge_send(wait:false)}) that reached a terminal phase
* (CB-588).
* Mechanism (b) of CB-307: a dedicated, status-gated push loop that nudges the primary's own
* herdr pane when a worker reply lands with no live {@code bridge_send} to resolve it.
*
* <p><strong>CB-590: one schedule per lead.</strong> Both kinds of work are triggered through
* their own entry point — {@link #onReplyQueued(String)} and
* {@link #onTicketTerminal(String, String, boolean)} — but both resolve the lead that should be
* nudged and coalesce onto a single per-lead reminder schedule, tracked in {@link #activeLeads}.
* Earlier this was two independent schedules (one keyed by worker target for replies, one keyed
* by lead for tickets) that could both decide to inject into the same pane in the same window —
* a race, not routine behaviour, but the expensive kind: it interrupts the lead's live turn
* twice. Collapsing to one schedule per lead makes that structurally impossible: at most one
* scheduled tick chain is ever live for a given lead (guarded by {@link #activeLeads}'
* compare-and-set), so at most one {@code agents.send} to that lead's pane is ever in flight.
* <p>The loop is triggered by {@link #onReplyQueued(String)} (called from
* {@link MessageService#reply} after the durable inbox publish). It checks four conditions
* at each tick via {@link #decide(String, int)}, then either injects a drain nudge,
* waits for the primary to become injectable, or stops reminding.
*
* <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, 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}.
* <p>Bounded: at most {@link #maxReminders} nudges per target, with a configurable backoff
* between them. The reply is never lost — the durable inbox is the backstop.
*
* <p><strong>CB-588 ticket nudges.</strong> {@link #onTicketTerminal(String, String, boolean)} is a
* second, independent entry point for an async delegation ticket ({@code bridge_send(wait:false)})
* reaching a terminal phase. That path takes the rendezvous fast path in {@link MessageService#reply}
* and never reaches {@link #onReplyQueued}, so without this the lead's own charter — prefer
* {@code wait:false} for anything non-trivial — was exactly the mode this loop failed to cover. It
* reuses the same status gating, bounded/backoff reminders, and metrics, keyed by the nudge-receiving
* lead terminal rather than the worker target so several tickets finishing together coalesce into one
* nudge. The two entry points do not interact: {@link #onReplyQueued} / {@link #decide} / their nudge
* text and bound are unchanged.
*/
public final class ReplyPushLoop {
private static final Logger log = LoggerFactory.getLogger(ReplyPushLoop.class);
static final String NUDGE_FORMAT = "Worker %s returned a reply — run bridge_poll(target=%s) to collect it";
/** Coalesced form, several uncollected replies for the same lead. */
static final String REPLIES_NUDGE_FORMAT =
"%d workers returned replies — run bridge_poll(target=...) for each to collect them: %s";
/** Singular form, one uncollected ticket. */
/** CB-588: singular form, one uncollected ticket. */
static final String TICKET_NUDGE_FORMAT =
"Ticket %s finished%s — run bridge_poll(ticket=%s) to collect it";
/** Coalesced form, several uncollected tickets for the same lead. */
/** CB-588: coalesced form, several uncollected tickets for the same lead. */
static final String TICKETS_NUDGE_FORMAT =
"%d tickets finished%s — run bridge_poll(ticket=...) for each to collect them: %s";
@@ -66,11 +56,11 @@ public final class ReplyPushLoop {
private final long backoffMs;
private final Metrics metrics; // CB-512: nullable — no registry in unit tests
/** Worker targets with a reply queued, and the lead to nudge about it, keyed by target. */
private final ConcurrentHashMap<String, String> pendingReplies = new ConcurrentHashMap<>();
/** Tickets that have gone terminal but not yet been polled, keyed by ticket. */
/** Track targets that have an active schedule. */
private final ConcurrentHashMap<String, Boolean> activeTargets = new ConcurrentHashMap<>();
/** CB-588: tickets that have gone terminal but not yet been polled, keyed by ticket. */
private final ConcurrentHashMap<String, PendingTicket> pendingTickets = new ConcurrentHashMap<>();
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets). */
/** CB-588: leads with an active ticket-reminder schedule. */
private final ConcurrentHashMap<String, Boolean> activeLeads = new ConcurrentHashMap<>();
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
@@ -99,79 +89,178 @@ public final class ReplyPushLoop {
}
}
// --- pending-work lookups (package-private for unit-testing) -------------------------------
// --- decision logic (package-private for unit-testing) -------------------------------------
/** The action the loop should take for a lead at the given reminder count. */
/** The action the loop should take for a target at the given reminder count. */
enum Action { INJECT, WAIT_BUSY, STOP }
/**
* Reply targets still pending for {@code lead} — registered via {@link #onReplyQueued} and
* whose inbox still holds an unacked message. A target whose inbox has since drained (acked,
* or collected via a live {@code bridge_send} rendezvous instead) is dropped from
* {@link #pendingReplies} here rather than lingering forever; there is no explicit "reply
* collected" callback the way {@link #ticketCollected} exists for tickets, so the inbox itself
* is the only signal.
* Pure decision function: examine the current state and return what the loop should do.
*
* @param target the worker session (target terminal id)
* @param reminderCount how many nudges have been sent so far for this target
* @return the action the caller should take
*/
private Set<String> pendingReplyTargetsFor(String lead) {
Set<String> result = new HashSet<>();
for (var entry : pendingReplies.entrySet()) {
String target = entry.getKey();
String owningLead = entry.getValue();
if (!lead.equals(owningLead)) continue;
if (inbox.peek(target).isEmpty()) {
pendingReplies.remove(target, owningLead);
continue;
}
result.add(target);
Action decide(String target, int reminderCount) {
// CB-532: the destination is per-delegation — the lead that sent this worker its work, not
// "the primary". With two leads orchestrating one fleet the singular question has no right
// answer, and answering it anyway interrupted whichever lead happened to call bridge_send
// first with results it never asked for.
var nudgeTarget = primaryRegistry.nudgeTargetFor(target);
if (nudgeTarget.isEmpty()) {
log.debug("push: no lead is known to be waiting on {}, stopping reminder", target);
return Action.STOP;
}
return result;
if (inbox.peek(target).isEmpty()) {
log.debug("push: inbox empty for {}, stopping reminder", target);
return Action.STOP;
}
if (reminderCount >= maxReminders) {
log.debug("push: reminder cap ({}) reached for {}, stopping", maxReminders, target);
countNudge("exhausted");
return Action.STOP;
}
String leadTerminal = nudgeTarget.get();
AgentStatus status;
try {
status = agents.status(leadTerminal);
} catch (RuntimeException e) {
log.debug("push: status check failed for lead {}, will retry", leadTerminal, e);
return Action.WAIT_BUSY;
}
if (status.injectable()) {
return Action.INJECT;
}
log.debug("push: lead {} is {} (not injectable), waiting", leadTerminal, status);
return Action.WAIT_BUSY;
}
// --- public entrypoint ---------------------------------------------------------------------
/**
* Called when a reply is queued for {@code target}. Idempotent per target: a second call while
* a schedule is active is a no-op. The schedule nudges the primary, then schedules a follow-up
* check (reminder on backoff, or re-check on WAIT_BUSY), until the inbox is empty or the cap
* is reached.
*/
public void onReplyQueued(String target) {
if (activeTargets.putIfAbsent(target, Boolean.TRUE) != null) {
log.debug("push: already active for {}, ignoring duplicate trigger", target);
return; // already scheduled
}
log.debug("push: starting reminder loop for {}", target);
scheduleNext(target, 0);
}
/** Execute one loop tick — called on the scheduler thread. */
private void tick(String target, int reminderCount) {
var action = decide(target, reminderCount);
switch (action) {
case INJECT -> {
injectNudge(target, reminderCount);
scheduleNext(target, reminderCount + 1);
}
// Re-check after the configured backoff; the primary may become injectable soon.
case WAIT_BUSY -> scheduleNext(target, reminderCount);
case STOP -> {
activeTargets.remove(target);
log.debug("push: reminder loop ended for {}", target);
}
}
}
/** Send the nudge and log the event. */
private void injectNudge(String target, int reminderCount) {
// Re-read rather than threading it down from decide(): the delegating lead can change
// between the decision and the injection, and the nudge should follow the current one.
var lead = primaryRegistry.nudgeTargetFor(target);
if (lead.isEmpty()) {
log.debug("push: lead for {} disappeared before the nudge could be sent", target);
return;
}
String leadTerminal = lead.get();
String nudge = NUDGE_FORMAT.formatted(target, target);
try {
agents.send(leadTerminal, nudge);
log.debug("push: nudge {}/{} sent to lead {} for target {}",
reminderCount + 1, maxReminders, leadTerminal, target);
countNudge("delivered");
} catch (RuntimeException e) {
log.warn("push: failed to nudge lead {} for target {} (reminder {}/{}): {}",
leadTerminal, target, reminderCount + 1, maxReminders, e.toString());
}
}
/** Schedule the next tick on the scheduler thread pool. */
private void scheduleNext(String target, int nextReminderCount) {
scheduler.schedule(() -> tick(target, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS);
}
// --- CB-588: async ticket terminal nudges ---------------------------------------------------
/** A ticket awaiting collection: which lead to nudge, and whether it ended in failure. */
private record PendingTicket(String ticket, String lead, boolean failed) {
}
/** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */
private List<PendingTicket> pendingTicketsFor(String lead) {
return pendingTickets.values().stream().filter(t -> lead.equals(t.lead())).toList();
}
/** Ticket ids still pending for {@code lead} — a plain snapshot for race comparison. */
private Set<String> pendingTicketIdsFor(String lead) {
return pendingTicketsFor(lead).stream().map(PendingTicket::ticket)
.collect(Collectors.toUnmodifiableSet());
/**
* Called when an async delegation ticket ({@code bridge_send(wait:false)}, CB-107) reaches a
* terminal phase — DONE or a failure. Unlike {@link #onReplyQueued}, which nudges about the
* durable-inbox no-waiter path, this covers the path {@code MessageService.reply} takes when a
* fire-and-poll send's own rendezvous waiter resolves the reply directly: that path returns
* before {@link #onReplyQueued} is ever called, so without this entry point a ticket finishing
* that way never nudged anyone (CB-588 / gitea #72).
*
* <p>Idempotent per lead: several tickets going terminal for the same lead while its schedule is
* already active coalesce onto that schedule's next tick rather than firing a nudge each.
*
* @param ticket the ticket to nudge about
* @param target the worker session the ticket was sent to — resolves which lead delegated it
* @param failed whether the ticket ended in a failure phase rather than {@code DONE}
*/
public void onTicketTerminal(String ticket, String target, boolean failed) {
var lead = primaryRegistry.nudgeTargetFor(target);
if (lead.isEmpty()) {
log.debug("push: no lead is known to be waiting on ticket {} (target {}), skipping nudge",
ticket, target);
return;
}
pendingTickets.put(ticket, new PendingTicket(ticket, lead.get(), failed));
if (activeLeads.putIfAbsent(lead.get(), Boolean.TRUE) != null) {
log.debug("push: ticket reminder loop already active for lead {}, {} coalesced in",
lead.get(), ticket);
return;
}
log.debug("push: starting ticket reminder loop for lead {}", lead.get());
scheduleTicketTick(lead.get(), 0);
}
/**
* Pure decision function: examine everything pending for {@code lead} — reply targets and
* tickets alike — and return what the loop should do.
*
* <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
* Called when a ticket's terminal state has been collected via {@code bridge_poll}. Removes it
* from the pending set so a scheduled tick — and any nudge it sends — never names a ticket the
* lead already has (CB-588 acceptance #5). A ticket that was never pending (unknown ticket, or
* one nudged with no push loop configured) is a no-op.
*/
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);
public void ticketCollected(String ticket) {
pendingTickets.remove(ticket);
}
/** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */
private List<PendingTicket> pendingFor(String lead) {
return pendingTickets.values().stream().filter(t -> lead.equals(t.lead())).toList();
}
/**
* Pure decision function for ticket nudges, mirroring {@link #decide(String, int)} but keyed by
* the nudge-receiving lead terminal rather than the worker session — several tickets from
* different workers delegated by the same lead coalesce onto it.
*/
Action decideTickets(String lead, int reminderCount) {
if (pendingFor(lead).isEmpty()) {
log.debug("push: nothing pending for lead {}, stopping ticket reminder", lead);
return Action.STOP;
}
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);
if (reminderCount >= maxReminders) {
log.debug("push: ticket reminder cap ({}) reached for lead {}, stopping", maxReminders, lead);
countNudge("exhausted");
return Action.STOP;
}
@@ -189,194 +278,95 @@ public final class ReplyPushLoop {
return Action.WAIT_BUSY;
}
// --- public entrypoints ----------------------------------------------------------------------
/**
* Called when a reply is queued for {@code target}. Resolves the lead delegating to
* {@code target} (CB-532) and coalesces onto that lead's single reminder schedule — starting
* one if none is active, joining an already-active one otherwise. A no-op if no lead is known
* to be waiting on {@code target}: there is nobody to nudge yet, and the durable inbox is the
* backstop until a lead is recorded.
*/
public void onReplyQueued(String target) {
var lead = primaryRegistry.nudgeTargetFor(target);
if (lead.isEmpty()) {
log.debug("push: no lead is known to be waiting on {}, skipping reminder", target);
return;
}
pendingReplies.put(target, lead.get());
startOrCoalesce(lead.get());
}
/**
* Called when an async delegation ticket ({@code bridge_send(wait:false)}, CB-107) reaches a
* terminal phase — DONE or a failure. Unlike {@link #onReplyQueued}, which nudges about the
* durable-inbox no-waiter path, this covers the path {@code MessageService.reply} takes when a
* fire-and-poll send's own rendezvous waiter resolves the reply directly: that path returns
* before {@link #onReplyQueued} is ever called, so without this entry point a ticket finishing
* that way never nudged anyone (CB-588 / gitea #72).
*
* <p>Resolves the delegating lead the same way {@link #onReplyQueued} does and coalesces onto
* the same per-lead schedule (CB-590) — several tickets, or a ticket and a reply, finishing
* for the same lead while its schedule is already active all ride the existing schedule's next
* tick rather than firing a nudge each.
*
* @param ticket the ticket to nudge about
* @param target the worker session the ticket was sent to — resolves which lead delegated it
* @param failed whether the ticket ended in a failure phase rather than {@code DONE}
*/
public void onTicketTerminal(String ticket, String target, boolean failed) {
var lead = primaryRegistry.nudgeTargetFor(target);
if (lead.isEmpty()) {
log.debug("push: no lead is known to be waiting on ticket {} (target {}), skipping nudge",
ticket, target);
return;
}
pendingTickets.put(ticket, new PendingTicket(ticket, lead.get(), failed));
startOrCoalesce(lead.get());
}
/**
* Called when a ticket's terminal state has been collected via {@code bridge_poll}. Removes it
* from the pending set so a scheduled tick — and any nudge it sends — never names a ticket the
* lead already has. A ticket that was never pending (unknown ticket, or one nudged with no push
* loop configured) is a no-op.
*/
public void ticketCollected(String ticket) {
pendingTickets.remove(ticket);
}
// --- the schedule ----------------------------------------------------------------------------
/** Start a reminder schedule for {@code lead}, or join the one already running. */
private void startOrCoalesce(String lead) {
if (activeLeads.putIfAbsent(lead, Boolean.TRUE) != null) {
log.debug("push: reminder loop already active for lead {}, work coalesced in", lead);
return;
}
log.debug("push: starting reminder loop for lead {}", lead);
scheduleNext(lead, 0, 0);
}
/** Execute one loop tick — called on the scheduler thread. */
private void tick(String lead, int replyReminderCount, int ticketReminderCount) {
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
Set<String> ticketsBefore = pendingTicketIdsFor(lead);
var action = decide(lead, replyReminderCount, ticketReminderCount);
/** Execute one ticket-loop tick — called on the scheduler thread. */
private void ticketTick(String lead, int reminderCount) {
Set<String> pendingBefore = pendingIdsFor(lead);
var action = decideTickets(lead, reminderCount);
switch (action) {
case INJECT -> {
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);
injectTicketNudge(lead, reminderCount);
scheduleTicketTick(lead, reminderCount + 1);
}
// Re-check after the configured backoff; the lead may become injectable soon.
case WAIT_BUSY -> scheduleNext(lead, replyReminderCount, ticketReminderCount);
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore);
case WAIT_BUSY -> scheduleTicketTick(lead, reminderCount);
case STOP -> stopOrRestartTicketLoop(lead, pendingBefore);
}
}
/** Ticket IDs pending for {@code lead} right now, as a plain snapshot for race comparison. */
private Set<String> pendingIdsFor(String lead) {
return pendingFor(lead).stream().map(PendingTicket::ticket).collect(Collectors.toUnmodifiableSet());
}
/**
* Release {@code lead}'s active-schedule slot, then restart it only if work landed that
* {@code repliesBefore} / {@code ticketsBefore} — the snapshots taken just before this tick's
* decision — did not already account for. {@link #onReplyQueued} / {@link #onTicketTerminal}
* read {@link #activeLeads} to decide whether to coalesce onto an existing schedule or start
* one, so work that lands between {@link #decide} returning {@link Action#STOP} and this
* removal running sees the (soon-to-be-stale) slot as occupied, coalesces onto a schedule that
* is about to die, and gets no nudge scheduled at all — a lost nudge, exactly what CB-588 (and
* now CB-590) exist to remove (originally found in review, gitea PR #73, for the ticket-only
* loop; carried forward here for the unified one).
* Release {@code lead}'s active-schedule slot, then restart it only if a ticket landed that
* {@code pendingBefore} — the snapshot taken just before this tick's decision — did not already
* account for. {@code onTicketTerminal} reads {@code activeLeads} to decide whether to coalesce
* onto an existing schedule or start one, so a ticket that lands between {@code decideTickets}
* returning {@link Action#STOP} and this removal running sees the (soon-to-be-stale) slot as
* occupied, coalesces onto a schedule that is about to die, and gets no nudge scheduled at all —
* a lost nudge, the exact failure CB-588 exists to remove (found in review, gitea PR #73).
*
* <p>Restarting on ANY non-empty pending set would be wrong: when STOP is reached because the
* reminder cap was hit rather than the backlog draining, the same never-collected work is
* expected to still be sitting there — that is the cap doing its job — and restarting would
* nudge about it forever, defeating the bound. Diffing the current pending sets against the
* "before" snapshots tells the two cases apart: an item present before this tick's decision is
* stale backlog, not a race; only an item absent from the "before" snapshot can only have
* arrived during the decision-to-release window, which is exactly the race this method closes.
* <p>Restarting on ANY non-empty {@code pendingFor(lead)} would be wrong: when STOP is reached
* because the reminder cap was hit rather than the backlog draining, the same never-collected
* ticket is expected to still be sitting there — that is the cap doing its job — and restarting
* would nudge about it forever, defeating the bound (the original CB-307 bounded-reminder
* guarantee, carried into CB-588 by acceptance criterion #7 — this exact regression showed up as
* two existing tests failing once a naive "any pending ticket restarts" version of this fix went
* in: {@code successfulTicketNudgeIncrementsDelivered} and {@code ticketNudgesSendUpToCapThenStop}).
* Diffing the current pending set against {@code pendingBefore} tells the two cases apart: a
* ticket present before this tick's decision is stale backlog, not a race; only a ticket absent
* from {@code pendingBefore} can only have arrived during the decision-to-release window, which is
* exactly the race this method closes.
*
* <p>Package-private so a test can drive the interleaving directly rather than trying to force a
* genuine thread race.
* genuine thread race: pass the exact {@code pendingBefore} snapshot a race requires (or does
* not) and call this to prove the recheck responds correctly either way.
*
* <p>Terminates rather than spinning: this method restarts the schedule at most once per call,
* and a fresh {@link #onReplyQueued} / {@link #onTicketTerminal} racing the recheck below still
* terminates in one of two ways — either it observes the slot already vacated (by the
* {@code activeLeads.remove} above, which happens-before this recheck in program order) and
* claims it itself, or it lands first and this recheck then observes its work in
* {@link #pendingReplies} / {@link #pendingTickets} and reclaims the slot instead. Exactly one
* side always wins; neither can miss the other, so this never loops on its own account.
* <p>Terminates rather than spinning: this method restarts the schedule at most once per call, and
* a fresh {@link #onTicketTerminal} racing the recheck below still terminates in one of two ways —
* either it observes the slot already vacated (by the {@code activeLeads.remove} above, which
* happens-before this recheck in program order) and claims it itself, or it lands first and this
* recheck then observes its ticket in {@code pendingTickets} and reclaims the slot instead. Exactly
* one side always wins; neither can miss the other, so this never loops on its own account.
*/
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore) {
void stopOrRestartTicketLoop(String lead, Set<String> pendingBefore) {
activeLeads.remove(lead);
boolean racedIn = pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t))
|| 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, 0);
boolean ticketRacedIn = pendingFor(lead).stream().anyMatch(t -> !pendingBefore.contains(t.ticket()));
if (ticketRacedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
log.debug("push: a ticket for lead {} raced the reminder loop's stop — restarting", lead);
scheduleTicketTick(lead, 0);
return;
}
log.debug("push: reminder loop ended for lead {}", lead);
log.debug("push: ticket reminder loop ended for lead {}", lead);
}
/** Send one combined nudge covering everything currently pending for {@code lead}. */
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);
List<PendingTicket> tickets = pendingTicketsFor(lead);
if (replyTargets.isEmpty() && tickets.isEmpty()) {
log.debug("push: pending work for lead {} drained before the nudge could be sent", lead);
/** Send the coalesced ticket nudge and log the event. */
private void injectTicketNudge(String lead, int reminderCount) {
// Re-read rather than threading it down from decideTickets(): a ticket can be collected (or
// another can arrive) between the decision and the injection.
List<PendingTicket> pending = pendingFor(lead);
if (pending.isEmpty()) {
log.debug("push: pending tickets for lead {} drained before the nudge could be sent", lead);
return;
}
String nudge = formatNudge(replyTargets, tickets);
String nudge = formatTicketsNudge(pending);
try {
agents.send(lead, nudge);
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());
log.debug("push: ticket nudge {}/{} sent to lead {} for {} ticket(s)",
reminderCount + 1, maxReminders, lead, pending.size());
countNudge("delivered");
} catch (RuntimeException e) {
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}): {}",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, e.toString());
log.warn("push: failed to nudge lead {} for {} ticket(s) (reminder {}/{}): {}",
lead, pending.size(), reminderCount + 1, maxReminders, e.toString());
}
}
/** Schedule the next tick on the scheduler thread pool. */
private void scheduleNext(String lead, int nextReplyReminderCount, int nextTicketReminderCount) {
scheduler.schedule(() -> tick(lead, nextReplyReminderCount, nextTicketReminderCount),
backoffMs, TimeUnit.MILLISECONDS);
/** Schedule the next ticket-loop tick on the scheduler thread pool. */
private void scheduleTicketTick(String lead, int nextReminderCount) {
scheduler.schedule(() -> ticketTick(lead, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS);
}
// --- nudge formatting ------------------------------------------------------------------------
/** Render everything pending for one lead as a single nudge line. */
private static String formatNudge(Set<String> replyTargets, List<PendingTicket> tickets) {
List<String> parts = new ArrayList<>();
if (!replyTargets.isEmpty()) {
parts.add(formatRepliesNudge(replyTargets));
}
if (!tickets.isEmpty()) {
parts.add(formatTicketsNudge(tickets));
}
return String.join(" | ", parts);
}
/** Render one or several pending reply targets. */
private static String formatRepliesNudge(Set<String> targets) {
if (targets.size() == 1) {
String target = targets.iterator().next();
return NUDGE_FORMAT.formatted(target, target);
}
String ids = String.join(", ", targets);
return REPLIES_NUDGE_FORMAT.formatted(targets.size(), ids);
}
/** Render one or several pending tickets. */
/** Render one or several pending tickets as a single nudge line. */
private static String formatTicketsNudge(List<PendingTicket> pending) {
if (pending.size() == 1) {
PendingTicket t = pending.get(0);
@@ -393,22 +383,23 @@ public final class ReplyPushLoop {
// --- lifecycle -----------------------------------------------------------------------------
/**
* Whether any reminder loop is currently active for some lead (CB-551). The idle-lead heartbeat
* uses this to stand aside: while the push loop is actively nudging a lead, a concurrent
* heartbeat injection would start a second competing turn in the same pane — racing loops
* multiply turns and context burn. "Active" means a schedule exists in {@link #activeLeads},
* which now covers both reply-queued (CB-307) and ticket-terminal (CB-588) work (CB-590) —
* bounded by what has been triggered, not by any persistent state.
* Whether any reminder loop is currently active for some target (CB-551). The idle-lead heartbeat
* uses this to stand aside: while the push loop is actively nudging the lead, a concurrent
* heartbeat injection would start a second competing turn in the same pane — racing loops multiply
* turns and context burn. "Active" means a schedule exists in {@link #activeTargets} or
* {@link #activeLeads} (CB-588 ticket nudges are a second source of pane injections the heartbeat
* must equally stand aside for); the sets are bounded by what has been triggered, not by any
* persistent state.
*/
public boolean isActive() {
return !activeLeads.isEmpty();
return !activeTargets.isEmpty() || !activeLeads.isEmpty();
}
/** Shut down the scheduler. Outstanding reminders are cancelled. */
public void stop() {
scheduler.shutdownNow();
activeTargets.clear();
activeLeads.clear();
pendingReplies.clear();
pendingTickets.clear();
}
@@ -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)
@@ -1,264 +0,0 @@
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 the (unfixed) unsynchronized sweep's iteration is a real, observable window
* a concurrently-started publish can land in — not just a best case, single-entry sprint. */
private static final int STALE_PUBLISHES = 100_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);
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) {
// resolved (failed by the sweep) — that is exactly what this thread is here for
}
});
}
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();
// A short, deliberate head start: with STALE_PUBLISHES this large, the (unfixed) sweep's own
// iteration takes several milliseconds, so this guarantees the sweep has already begun —
// and, once guarded, is already holding publishChannelLock — before "fresh" attempts to
// register. Without this head start, "fresh" sometimes wins the race for the lock and
// registers before the sweep even starts, which is the accepted "already in flight when
// recovery fires" case (correctly failed either way) rather than the bug under test.
Thread.sleep(5);
// 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 20,000 stale threads' own lock acquisitions) 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;
}
}
@@ -20,7 +20,6 @@ import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import static org.junit.jupiter.api.Assertions.*;
@@ -28,12 +27,6 @@ import static org.junit.jupiter.api.Assertions.*;
* Unit tests for {@link ReplyPushLoop}: decision logic, nudge injection, idempotency,
* bounded reminders, and stop conditions.
*
* <p>CB-590 collapsed the CB-307 reply-nudge schedule and the CB-588 ticket-nudge schedule into
* one schedule per lead ({@link ReplyPushLoop#decide}), so most tests below register pending work
* through the public entry points ({@code onReplyQueued} / {@code onTicketTerminal}) before
* exercising {@code decide} directly, mirroring how the two entry points now share one decision
* function keyed by the lead terminal rather than by worker target.
*
* <p>Uses a {@link RecordingHerdrClient} that synchronizes access to its call list so the
* scheduler thread and test thread never have memory ordering issues. The {@code decide()}
* tests use a simple client with no concurrency concern.
@@ -65,50 +58,38 @@ class ReplyPushLoopTest {
// --- decide() logic ------------------------------------------------------------------------
@Test
void onReplyQueuedWithNoKnownLeadNeverStartsASchedule() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 50);
loop.onReplyQueued(WORKER); // no lead known -> never registered, never scheduled
assertFalse(loop.isActive(), "no lead known means nothing to nudge yet");
Thread.sleep(150);
assertEquals(0, rec.sendCount(), "must not nudge when no lead is known to be waiting");
void decideWithoutPrimaryIsStop() {
agents = agentWithStatus("idle");
var loop = new ReplyPushLoop(
new PrimaryRegistry(null), agents, inbox, scheduler, 5, 100);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(WORKER, 0));
}
@Test
void decideWithNothingPendingIsStop() {
void decideWithEmptyInboxIsStop() {
agents = agentWithStatus("idle");
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
}
@Test
void decideAtCapIsStop() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
var loop = loop(2, 100);
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100).decide(WORKER, 2));
}
@Test
void decideUnderCapWithInjectablePrimaryIsInject() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
}
@Test
void decideUnderCapWithBlockedPrimaryIsInject() {
agents = agentWithStatus("blocked");
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
"BLOCKED is injectable");
}
@@ -116,9 +97,7 @@ class ReplyPushLoopTest {
void decideUnderCapWithDonePrimaryIsInject() {
agents = agentWithStatus("done");
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
"DONE is injectable");
}
@@ -126,29 +105,23 @@ class ReplyPushLoopTest {
void decideUnderCapWithBusyPrimaryIsWaitBusy() {
agents = agentWithStatus("working");
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
}
@Test
void decideUnderCapWithUnknownPrimaryIsWaitBusy() {
agents = agentWithStatus("unknown");
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
}
@Test
void decideStopsAfterInboxIsEmptied() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
inbox.ack(WORKER, "m1");
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
}
// --- onReplyQueued integration -------------------------------------------------------------
@@ -228,19 +201,12 @@ class ReplyPushLoopTest {
assertTrue(nudge.contains("bridge_poll(target=term_worker)"));
}
@Test
void repliesNudgeFormatIsCorrect() {
String multi = ReplyPushLoop.REPLIES_NUDGE_FORMAT.formatted(2, "term_worker1, term_worker2");
assertTrue(multi.contains("2 workers"));
assertTrue(multi.contains("bridge_poll(target=...)"));
}
// --- CB-588: async ticket terminal nudges — decide() logic on tickets -----------------------
// --- CB-588: async ticket terminal nudges — decideTickets() logic --------------------------
@Test
void decideTicketsWithNothingPendingIsStop() {
agents = agentWithStatus("idle");
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop().decideTickets(PRIMARY, 0));
}
@Test
@@ -248,7 +214,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, 0, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2));
}
@Test
@@ -256,7 +222,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, 0));
assertEquals(ReplyPushLoop.Action.INJECT, loop.decideTickets(PRIMARY, 0));
}
@Test
@@ -264,7 +230,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, 0));
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decideTickets(PRIMARY, 0));
}
@Test
@@ -272,7 +238,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, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 0));
}
// --- CB-588: onTicketTerminal integration ---------------------------------------------------
@@ -289,7 +255,7 @@ class ReplyPushLoopTest {
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("task-1"), "nudge should name the ticket");
assertTrue(nudge.contains("bridge_poll(ticket="), "nudge should name the exact ticket-poll call");
assertFalse(nudge.contains("bridge_poll(target="), "a ticket-only nudge must not tell the lead to run the target-poll call");
assertFalse(nudge.contains("bridge_poll(target="), "a ticket nudge must not tell the lead to run the target-poll call");
}
@Test
@@ -380,11 +346,11 @@ class ReplyPushLoopTest {
void aTicketStillPendingWhenTheLoopStopsIsNotStranded() {
// Regression for the race a reviewer found in gitea PR #73: onTicketTerminal's
// activeLeads.putIfAbsent can see the lead's slot as still occupied a moment before
// decide's STOP releases it, so the ticket coalesces onto a schedule that is about to
// decideTickets' STOP releases it, so the ticket coalesces onto a schedule that is about to
// die and nothing ever nudges about it. Forcing that exact thread interleaving is not
// reliable, so this drives stopOrRestart — the STOP path's own release-and-recheck —
// reliable, so this drives stopOrRestartTicketLoop — the STOP path's own release-and-recheck —
// directly, arranging the state it must not lose a ticket in: a ticket pending for the lead
// that was NOT part of the pre-decision snapshot (ticketsBefore=empty), standing in for one
// that was NOT part of the pre-decision snapshot (pendingBefore=empty), standing in for one
// that races in during the decision-to-release window.
var rec = recordingClient();
agents = new AgentControl(rec);
@@ -392,9 +358,9 @@ class ReplyPushLoopTest {
loop.onTicketTerminal("task-1", WORKER, false); // pendingTickets={task-1}; activeLeads={PRIMARY}
// Stand in for the scheduler thread reaching decide==STOP for this lead — with nothing
// Stand in for the scheduler thread reaching decideTickets==STOP for this lead — with nothing
// pending at decide time — while task-1 races in before the release below runs.
loop.stopOrRestart(PRIMARY, Set.of(), Set.of());
loop.stopOrRestartTicketLoop(PRIMARY, Set.of());
assertTrue(loop.isActive(), "a ticket that raced the loop's stop must reclaim the schedule "
+ "slot, not be stranded with no schedule left to ever nudge about it");
@@ -402,67 +368,23 @@ class ReplyPushLoopTest {
@Test
void aStaleUncollectedTicketAtCapDoesNotRestartTheLoop() {
// The other direction of the same fix: restarting on ANY non-empty pending set would be
// The other direction of the same fix: restarting on ANY non-empty pendingFor(lead) would be
// wrong. When STOP is reached because the reminder cap was hit, the same never-collected
// ticket is expected to still be there — that is the cap doing its job (acceptance criterion
// #5: no spin / nudges stay bounded). task-1 here was already accounted for at decide time
// (it is in ticketsBefore), so it must not restart the loop just because it is still sitting
// there.
// #7: nudges stay bounded). task-1 here was already accounted for at decide time (it is in
// pendingBefore), so it must not restart the loop just because it is still sitting there.
var rec = recordingClient();
agents = new AgentControl(rec);
ReplyPushLoop loop = loop(1, 100_000);
loop.onTicketTerminal("task-1", WORKER, false); // pendingTickets={task-1}; activeLeads={PRIMARY}
loop.stopOrRestart(PRIMARY, Set.of(), Set.of("task-1"));
loop.stopOrRestartTicketLoop(PRIMARY, Set.of("task-1"));
assertFalse(loop.isActive(), "a stale ticket already accounted for at decide time must not "
+ "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");
@@ -480,103 +402,7 @@ class ReplyPushLoopTest {
void inboxNudgeStillUsesTheOriginalTargetPollCall() {
String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER);
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
"CB-588/CB-590 must not change the CB-307 inbox nudge's call shape");
}
// --- CB-590: one schedule per lead — no overlap, no lost nudges -----------------------------
@Test
void replyAndTicketForTheSameLeadCoalesceIntoOneSendNeverOverlapping() throws Exception {
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
var loop = loop(1, 300); // backoff wide enough that both entry points land before the first tick
loop.onReplyQueued(WORKER);
loop.onTicketTerminal("task-1", WORKER, false);
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one combined nudge should have been sent");
Thread.sleep(300);
assertEquals(1, rec.sendCount(),
"a reply and a ticket for the same lead must coalesce onto ONE schedule — "
+ "two nudge injections into the same lead pane must never overlap");
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
"the combined nudge must still mention the reply: " + nudge);
assertTrue(nudge.contains("task-1"), "the combined nudge must still mention the ticket: " + nudge);
}
@Test
void aReplyQueuedWhileTheLeadIsBusyIsNotLostWhenATicketArrivesToo() throws Exception {
// The lead is busy for its first two status checks, then becomes injectable. A reply is
// queued while busy; a ticket for the same lead arrives before the lead frees up. Neither
// may be dropped — deferred is fine, lost is not (acceptance criterion #2).
var rec = new BusyThenIdleHerdrClient(2);
agents = new AgentControl(rec);
inbox.publish(WORKER, "m1", "hello");
var loop = loop(1, 50); // cap=1: WAIT_BUSY doesn't count against it, so exactly one send once injectable
loop.onReplyQueued(WORKER); // schedule starts, first tick(s) WAIT_BUSY
loop.onTicketTerminal("task-1", WORKER, false); // coalesces onto the same waiting schedule
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
"once the lead becomes injectable, the deferred work must still be nudged");
Thread.sleep(200);
assertEquals(1, rec.sendCount(), "exactly one nudge once injectable — reply and ticket coalesced");
String nudge = rec.sentParams().getFirst().getValue().toString();
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"), "the reply must not be dropped: " + nudge);
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");
"CB-588 must not change the CB-307 inbox nudge's call shape");
}
// --- metrics (CB-512) ----------------------------------------------------------------------
@@ -604,10 +430,8 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
Metrics metrics = new Metrics();
var loop = loop(2, 100, metrics);
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100, metrics).decide(WORKER, 2));
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the reminder cap must count as exhausted");
@@ -635,7 +459,7 @@ class ReplyPushLoopTest {
var loop = loop(2, 100_000, metrics);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2));
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the ticket reminder cap must count as exhausted");
@@ -723,49 +547,4 @@ class ReplyPushLoopTest {
private static RecordingHerdrClient recordingClient() {
return new RecordingHerdrClient();
}
/**
* Thread-safe fake that reports {@code working} (not injectable) for its first
* {@code busyChecks} status calls, then {@code idle} forever after — used to prove work queued
* while the lead is busy is deferred, not dropped, once it becomes injectable.
*/
private static final class BusyThenIdleHerdrClient implements HerdrClient {
private final List<Map.Entry<String, Object>> calls =
Collections.synchronizedList(new ArrayList<>());
private final AtomicInteger statusChecks = new AtomicInteger();
private final int busyChecks;
volatile CountDownLatch sendLatch = new CountDownLatch(1);
BusyThenIdleHerdrClient(int busyChecks) {
this.busyChecks = busyChecks;
}
@Override
public JsonNode call(String method, Object params) {
if ("agent.get".equals(method)) {
String status = statusChecks.getAndIncrement() < busyChecks ? "working" : "idle";
return MAPPER.createObjectNode()
.set("agent", MAPPER.createObjectNode()
.put("terminal_id", PRIMARY)
.put("agent_status", status));
}
if ("agent.prompt".equals(method)) {
calls.add(Map.entry(method, params));
sendLatch.countDown();
}
return MAPPER.createObjectNode();
}
long sendCount() {
return calls.size();
}
List<Map.Entry<String, Object>> sentParams() {
return List.copyOf(calls);
}
@Override
public void close() {
}
}
}
+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)"