Compare commits

...

8 Commits

Author SHA1 Message Date
Dai Ha 8a837a2830 CB-599: surface a capacity refusal's reason instead of a bare 500
CI / contract (pull_request) Successful in 50s
CI / build (pull_request) Successful in 1m35s
PlacementException extends IllegalStateException, which neither BridgedApp
nor BridgeMcp's spawn catch blocks handled, so a maxLoad/quarantine/
all-exhausted refusal fell through to a blank 500 on REST and lost its
message on MCP. Catch it on both surfaces, ahead of the unrelated
IllegalArgumentException(unknown_profile) mapping, and return its message
structured: REST as {"error":"no_capacity","detail":...} with status 503,
MCP as an isError result prefixed "no capacity: ...".
2026-08-16 17:43:07 +02:00
ltms 613ece92dc Merge CB-594: make supervision and a working fleet possible at the same time
CI / contract (push) Successful in 49s
CI / build (push) Failing after 1m45s
The launchd unit was a CHANGEME template that had never been installed, and it could not have worked if it were: launchd does not source a login shell, so the daemon would have started with no forge or gateway token, and the failure would only appear much later as workers unable to open a PR.

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

AmqpReplyInboxRecoveryRaceTest drives the sweep and a real publish() against
each other directly (no live broker reconnect) using Proxy-backed fake AMQP
channels and a large in-flight backlog to make the race window observable;
confirmed it fails without the guard (reply "fresh" wrongly failed as
"connection recovered mid-publish") and passes with it.
2026-08-16 17:24:43 +02:00
Dai Ha 78ca24dc3f CB-590: collapse the CB-307 and CB-588 nudge schedules into one per lead
CI / build (pull_request) Successful in 1m0s
CI / contract (pull_request) Successful in 1m33s
Both reply-queued and ticket-terminal nudges could independently decide
to inject into the same lead pane in the same window, since they ran as
two separate schedules keyed differently (worker target vs. lead) that
never checked each other. Replace both with a single per-lead schedule
(activeLeads) that drains pending reply targets and pending tickets
together, sends at most one combined nudge per tick, and shares one
reminder cap across both sources — so two injections into the same pane
can no longer overlap, and work queued while the lead is busy is never
lost, only deferred.
2026-08-16 16:56:04 +02:00
Dai Ha 4fa6553db5 CB-527/CB-528: bound AMQP prefetch and confirm publishes before claiming durable
CI / build (pull_request) Successful in 1m3s
CI / contract (pull_request) Successful in 1m15s
CB-527: basicQos(prefetch) on the consume channel before basicConsume, configurable
via broker.prefetch (default 32), so an undrained inbox backlog stays on the broker
instead of growing the JVM heap without limit.

CB-528: publish moves to its own confirm-mode channel with mandatory=true and a
return listener, so an unroutable or unconfirmed reply now throws instead of
vanishing silently. The confirm callback checks the per-message returned flag
(set by the return listener, which the broker always fires before the matching
confirm) so an acked-but-returned publish is still reported as a failure. The ack
path stays on its own channel/lock and never waits on a publish confirm.
2026-08-16 16:55:40 +02:00
12 changed files with 1180 additions and 298 deletions
+6 -2
View File
@@ -434,10 +434,14 @@ guard:
# on a durable per-target queue (agent.<target>.inbox) and survive a restart — the broker
# redelivers anything the primary had not yet drained. Production default is LavinMQ; a stock
# RabbitMQ speaks the same AMQP 0-9-1, so it is a URI-only swap.
# uri → AMQP connection URI. No trailing slash ⇒ the default vhost "/"; an empty path ("/")
# is vhost "" and will NOT connect. Encode a named vhost as .../%2Fmyvhost.
# uri → AMQP connection URI. No trailing slash ⇒ the default vhost "/"; an empty path ("/")
# is vhost "" and will NOT connect. Encode a named vhost as .../%2Fmyvhost.
# prefetch → CB-527: consumer basicQos, capping how many unacked messages the inbox holds
# in-heap per owned target (the rest sits on the broker's durable queue instead of
# growing the JVM heap). Default 32 when omitted.
# broker:
# uri: amqp://guest:guest@127.0.0.1:5672
# prefetch: 32
# Active push-to-primary (CB-307 Stage 3). When a worker reply lands with no open bridge_send,
# the ReplyPushLoop injects a *drain nudge* (never the payload) into the primary's own herdr
@@ -352,8 +352,9 @@ public final class Bridged {
// connection, so keep the reference to close it in the ordered shutdown hook.
final ReplyInbox replyInbox;
if (cfg.broker() != null && cfg.broker().isConfigured()) {
replyInbox = AmqpReplyInbox.open(cfg.broker().uri());
log.info("reply inbox: AMQP broker (durable) at {}", cfg.broker().uri());
replyInbox = AmqpReplyInbox.open(cfg.broker().uri(), cfg.broker().prefetchOrDefault());
log.info("reply inbox: AMQP broker (durable) at {} (prefetch={})",
cfg.broker().uri(), cfg.broker().prefetchOrDefault());
} else {
replyInbox = new InMemoryReplyInbox();
log.info("reply inbox: in-memory (soft-state)");
@@ -5,6 +5,7 @@ import com.fasterxml.jackson.core.JsonParser;
import com.fasterxml.jackson.core.JsonToken;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
import dev.ltms.bridged.msg.AmqpReplyInbox;
import dev.ltms.bridged.peer.MemberRole;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -532,16 +533,24 @@ public record BridgedConfig(
* stays soft-state. Production default is LavinMQ; a stock RabbitMQ speaks the same AMQP 0-9-1
* and is a URI-only swap.
*
* @param uri AMQP connection URI, e.g. {@code amqp://guest:guest@127.0.0.1:5672/}. Blank/{@code null}
* ⇒ the broker block is treated as absent (in-memory adapter).
* @param uri AMQP connection URI, e.g. {@code amqp://guest:guest@127.0.0.1:5672/}. Blank/
* {@code null} ⇒ the broker block is treated as absent (in-memory adapter).
* @param prefetch CB-527: the consumer's {@code basicQos} prefetch count, bounding how many
* unacked messages the AMQP inbox holds in-heap per owned target. {@code null}/
* non-positive ⇒ {@link AmqpReplyInbox#DEFAULT_PREFETCH}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Broker(String uri) {
public record Broker(String uri, Integer prefetch) {
/** True when a usable broker URI is configured (an empty block does not enable AMQP). */
public boolean isConfigured() {
return uri != null && !uri.isBlank();
}
/** The prefetch to use, defaulting to {@link AmqpReplyInbox#DEFAULT_PREFETCH} when unset. */
public int prefetchOrDefault() {
return (prefetch != null && prefetch > 0) ? prefetch : AmqpReplyInbox.DEFAULT_PREFETCH;
}
}
/**
@@ -15,6 +15,7 @@ import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.placement.BackendQuarantine;
import dev.ltms.bridged.placement.PlacementException;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.session.WorktreeRequest;
@@ -692,6 +693,10 @@ public final class BridgeMcp {
return text(json(memberView(member)));
} catch (GuardException e) {
return error("subscription boundary: " + e.getMessage());
} catch (PlacementException e) {
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — distinct
// from "profile does not exist" below.
return error("no capacity: " + e.getMessage());
} catch (IllegalArgumentException e) {
return error(e.getMessage()); // unknown / no-default profile, or a refused resumeSessionId
} catch (PeerUnreachableException e) {
@@ -7,6 +7,7 @@ import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import com.rabbitmq.client.Recoverable;
import com.rabbitmq.client.RecoveryListener;
import com.rabbitmq.client.Return;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -14,7 +15,13 @@ import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.NavigableMap;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentSkipListMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
/**
* AMQP-backed {@link ReplyInbox} (CB-307 Stage 2): genuine cross-restart durability behind the same
@@ -29,6 +36,27 @@ import java.util.concurrent.ConcurrentHashMap;
* leaves them on the broker — it redelivers on reconnect. That is the durability the in-memory
* adapter cannot give, with the port contract preserved.
*
* <p><strong>Prefetch bounds the held backlog (CB-527).</strong> The consumer channel calls
* {@code basicQos} with a configurable prefetch count ({@link #DEFAULT_PREFETCH} unless the caller
* passes another value to {@link #open(String, int)}) before starting any consumer. Without a bound,
* the broker pushes its entire queue into {@link #held} the instant a target is {@link #own owned},
* so an undrained primary grows the JVM heap without limit and any queue-level control
* ({@code x-max-length}, per-message TTL) never fires because the queue never actually holds a
* backlog. Prefetch keeps the backlog where it is visible — on the broker — until the owner drains it.
*
* <p><strong>Publishes require a confirmed, routable delivery (CB-528).</strong> {@link #publish}
* runs on a channel separate from the consume/ack channel ({@link #channel}), so a slow or blocked
* publish confirm can never hold {@link #channelLock} and stall an ack — the ack path never waits on
* a publish confirm. That publish channel is in publisher-confirm mode and every publish sets the
* {@code mandatory} flag, so an unroutable publish (queue not declared, e.g. the owner never called
* {@link #own}) is returned by the broker instead of silently dropped. The broker sends the
* <em>return</em> for an unroutable message before the <em>confirm</em> that covers it — the ack/nack
* callback checks the returned-set at confirm time rather than assuming an ack means routed — so
* "confirmed" here means "durably queued", not merely "accepted by the broker". A returned or nacked
* (or un-confirmed within the timeout) publish surfaces as an {@link IllegalStateException} on the
* caller's thread; the caller — {@link MessageService#reply} — must not report success for a
* black-holed reply.
*
* <p><strong>Ownership is explicit.</strong> {@link #own} declares the queue and starts the consumer;
* {@link #release} cancels it. {@link #publish} sends to the queue but does <em>not</em> imply ownership
* and does not attach a consumer. This split is required by CB-308 federation, where one gateway may
@@ -54,6 +82,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
private static final String QUEUE_PREFIX = "agent.";
private static final String QUEUE_SUFFIX = ".inbox";
/** CB-527: the prefetch used when a caller does not pass an explicit value to {@link #open(String, int)}. */
public static final int DEFAULT_PREFETCH = 32;
/** How long {@link #publish} waits for its publisher confirm before failing the call (CB-528). */
private static final long CONFIRM_TIMEOUT_MS = 10_000L;
private final Connection connection;
private final Channel channel;
/** All channel operations (publish/declare/ack/cancel) serialize on this — a Channel is not thread-safe. */
@@ -63,39 +97,87 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
/** Targets whose queue is declared and consumer is running, mapped to their broker consumer tag. */
private final ConcurrentHashMap<String, String> consumerTags = new ConcurrentHashMap<>();
/**
* CB-528: a dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume
* + ack) so a publish confirm round trip never blocks under {@link #channelLock} and stalls an ack.
*/
private final Channel publishChannel;
private final Object publishChannelLock = new Object();
/** In-flight publishes awaiting their confirm, keyed by the publish channel's sequence number. */
private final ConcurrentSkipListMap<Long, Pending> pendingBySeq = new ConcurrentSkipListMap<>();
/**
* The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery
* tag. Assumes {@code msgId} is unique per in-flight publish: a second {@link #publish} for a
* {@code msgId} still awaiting its confirm would overwrite this entry and misdirect
* {@link #onReturn}'s lookup. Not reachable today — {@code MessageService.reply} generates a fresh
* {@code UUID} per call — so no guard is added for it.
*/
private final ConcurrentHashMap<String, Pending> pendingByMsgId = new ConcurrentHashMap<>();
/** A message pulled off the broker but not yet acked: its delivery-tag plus the port payload. */
private record Held(long deliveryTag, InboxMessage message) {}
/** Connect to {@code uri} (e.g. {@code amqp://guest:guest@127.0.0.1:5672/}) and open the inbox. */
/** A publish awaiting its confirm; {@link #returned} records whether the broker already returned it. */
private static final class Pending {
final String msgId;
final CompletableFuture<Void> confirmed = new CompletableFuture<>();
volatile boolean returned;
Pending(String msgId) {
this.msgId = msgId;
}
}
/** Connect to {@code uri} (e.g. {@code amqp://guest:guest@127.0.0.1:5672/}) with {@link #DEFAULT_PREFETCH}. */
public static AmqpReplyInbox open(String uri) {
return open(uri, DEFAULT_PREFETCH);
}
/** As {@link #open(String)}, with an explicit consumer prefetch (CB-527: caps the held backlog per target). */
public static AmqpReplyInbox open(String uri, int prefetch) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
return new AmqpReplyInbox(factory.newConnection("bridged-reply-inbox"));
return new AmqpReplyInbox(factory.newConnection("bridged-reply-inbox"), prefetch);
} catch (Exception e) {
throw new IllegalStateException("cannot connect to AMQP broker at " + uri, e);
}
}
/** Wrap an already-open connection (injection seam for the contract test). */
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for the contract test). */
AmqpReplyInbox(Connection connection) {
this(connection, DEFAULT_PREFETCH);
}
/** As above, with an explicit prefetch (injection seam for the contract test). */
AmqpReplyInbox(Connection connection, int prefetch) {
this.connection = connection;
try {
this.channel = connection.createChannel();
// CB-527: bound the held backlog per owned target — must be set before any own()/basicConsume.
this.channel.basicQos(prefetch);
this.publishChannel = connection.createChannel();
this.publishChannel.confirmSelect();
this.publishChannel.addReturnListener(this::onReturn);
this.publishChannel.addConfirmListener(this::onAck, this::onNack);
} catch (IOException e) {
throw new IllegalStateException("cannot open AMQP channel", e);
}
// On automatic recovery the broker redelivers unacked messages with FRESH delivery-tags; the
// tags we were holding are now stale. Drop the held snapshot so the re-attached consumer
// repopulates it with valid tags (dedup by msgId still prevents any double-queue).
// repopulates it with valid tags (dedup by msgId still prevents any double-queue). Any publish
// confirm still in flight when the connection dropped is equally stale — its sequence number
// meant nothing on the old channel and means nothing on the recovered one, so fail it now
// rather than let it silently ride out CONFIRM_TIMEOUT_MS.
if (connection instanceof Recoverable recoverable) {
recoverable.addRecoveryListener(new RecoveryListener() {
@Override
public void handleRecovery(Recoverable recoverable) {
held.clear();
failPendingPublishesOnRecovery();
log.info("AMQP connection recovered; cleared held replies for fresh redelivery");
}
@@ -141,6 +223,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
}
/**
* Publish {@code content} and block until the broker's publisher confirm for it lands (CB-528).
* Throws {@link IllegalStateException} if the message is returned as unroutable, nacked, or not
* confirmed within {@link #CONFIRM_TIMEOUT_MS} — the caller must treat that as a failed publish,
* not a lost-and-forgotten one.
*/
@Override
public void publish(String target, String msgId, String content) {
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
@@ -148,12 +236,35 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
.deliveryMode(2) // persistent — survives a broker restart
.contentType("text/plain")
.build();
try {
synchronized (channelLock) {
channel.basicPublish("", queueName(target), props, content.getBytes(StandardCharsets.UTF_8));
Pending pending = new Pending(msgId);
long seq;
synchronized (publishChannelLock) {
seq = publishChannel.getNextPublishSeqNo();
pendingBySeq.put(seq, pending);
pendingByMsgId.put(msgId, pending);
try {
publishChannel.basicPublish("", queueName(target), true, props,
content.getBytes(StandardCharsets.UTF_8));
} catch (IOException e) {
pendingBySeq.remove(seq, pending);
pendingByMsgId.remove(msgId, pending);
throw new IllegalStateException("cannot publish reply to " + queueName(target), e);
}
} catch (IOException e) {
throw new IllegalStateException("cannot publish reply to " + queueName(target), e);
}
try {
pending.confirmed.get(CONFIRM_TIMEOUT_MS, TimeUnit.MILLISECONDS);
} catch (ExecutionException e) {
Throwable cause = e.getCause();
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
} catch (TimeoutException e) {
throw new IllegalStateException("publish confirm for reply " + msgId + " to " + queueName(target)
+ " timed out after " + CONFIRM_TIMEOUT_MS + "ms — broker may be unreachable or overloaded", e);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new IllegalStateException("interrupted awaiting publish confirm for " + msgId, e);
} finally {
pendingBySeq.remove(seq, pending);
pendingByMsgId.remove(msgId, pending);
}
}
@@ -222,17 +333,115 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
};
}
/** Broker return for an unroutable {@code mandatory} publish — arrives BEFORE its confirm (CB-528). */
private void onReturn(Return r) {
String msgId = r.getProperties() == null ? null : r.getProperties().getMessageId();
Pending pending = msgId == null ? null : pendingByMsgId.get(msgId);
if (pending != null) {
pending.returned = true;
} else {
log.warn("AMQP return for reply {} (routingKey={}, {} {}) with no matching in-flight publish"
+ " — already resolved by a prior confirm", msgId, r.getRoutingKey(), r.getReplyCode(),
r.getReplyText());
}
}
private void onAck(long seq, boolean multiple) {
resolveConfirm(seq, multiple, true);
}
private void onNack(long seq, boolean multiple) {
resolveConfirm(seq, multiple, false);
}
/**
* Resolve every pending publish covered by this confirm (a single seq, or — {@code multiple} —
* every seq up to and including it). Checks {@link Pending#returned} at confirm time: since the
* broker's return for an unroutable message always precedes its confirm, an ack that arrives after
* a return means "confirmed but never routed", not "durably queued".
*/
private void resolveConfirm(long seq, boolean multiple, boolean ack) {
NavigableMap<Long, Pending> covered = multiple
? pendingBySeq.headMap(seq, true)
: pendingBySeq.subMap(seq, true, seq, true);
for (var it = covered.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
if (ack && !pending.returned) {
pending.confirmed.complete(null);
} else if (ack) {
pending.confirmed.completeExceptionally(new IllegalStateException(
"reply " + pending.msgId + " was returned as unroutable (queue not declared/owned)"));
} else {
pending.confirmed.completeExceptionally(new IllegalStateException(
"broker nacked publish of reply " + pending.msgId));
}
}
}
/**
* Fail every publish still awaiting its confirm — their sequence numbers are stale after recovery.
* Guarded by {@link #publishChannelLock}, the same lock {@link #publish} holds while it takes its
* sequence number and registers its {@link Pending}: without it, a {@link #publish} that starts
* after the connection has already recovered (so it publishes — and will be confirmed — on the
* <em>new</em> channel) can register between this sweep's iteration and its clear, and this sweep
* then fails a publish that actually succeeded. {@link #publish} only holds the lock for the
* seq/map-put/{@code basicPublish} — it awaits the confirm outside it — so this sweep can only ever
* wait for an in-flight {@code basicPublish} call to return, never for a broker round trip. No
* deadlock.
*
* <p>Package-private (rather than {@code private}) only so the unit test can drive it directly
* against a concurrent {@link #publish} without a live broker reconnect.
*/
void failPendingPublishesOnRecovery() {
synchronized (publishChannelLock) {
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
pending.confirmed.completeExceptionally(new IllegalStateException(
"AMQP connection recovered mid-publish; confirm status of reply " + pending.msgId
+ " is unknown"));
}
}
}
/**
* Fail every publish still awaiting its confirm with a clear, immediate error instead of leaving it
* to time out after {@link #CONFIRM_TIMEOUT_MS} once the channels are closed underneath it. Guarded
* by {@link #publishChannelLock} for the same reason as {@link #failPendingPublishesOnRecovery}.
*/
private void failPendingPublishesOnClose() {
synchronized (publishChannelLock) {
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
pending.confirmed.completeExceptionally(new IllegalStateException(
"AMQP reply inbox closed while publish of reply " + pending.msgId
+ " was still awaiting its confirm"));
}
}
}
private static String queueName(String target) {
return QUEUE_PREFIX + target + QUEUE_SUFFIX;
}
@Override
public void close() {
failPendingPublishesOnClose();
try {
channel.close();
} catch (Exception e) {
log.debug("AMQP channel close: {}", e.toString());
}
try {
publishChannel.close();
} catch (Exception e) {
log.debug("AMQP publish channel close: {}", e.toString());
}
try {
connection.close();
} catch (Exception e) {
@@ -8,6 +8,8 @@ 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;
@@ -16,35 +18,43 @@ import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
/**
* 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.
* 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).
*
* <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><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>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.
* <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}.
*/
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";
/** CB-588: singular form, one uncollected ticket. */
/** 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. */
static final String TICKET_NUDGE_FORMAT =
"Ticket %s finished%s — run bridge_poll(ticket=%s) to collect it";
/** CB-588: coalesced form, several uncollected tickets for the same lead. */
/** 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";
@@ -56,11 +66,11 @@ public final class ReplyPushLoop {
private final long backoffMs;
private final Metrics metrics; // CB-512: nullable — no registry in unit tests
/** 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. */
/** 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. */
private final ConcurrentHashMap<String, PendingTicket> pendingTickets = new ConcurrentHashMap<>();
/** CB-588: leads with an active ticket-reminder schedule. */
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets). */
private final ConcurrentHashMap<String, Boolean> activeLeads = new ConcurrentHashMap<>();
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
@@ -89,178 +99,79 @@ public final class ReplyPushLoop {
}
}
// --- decision logic (package-private for unit-testing) -------------------------------------
// --- pending-work lookups (package-private for unit-testing) -------------------------------
/** The action the loop should take for a target at the given reminder count. */
/** The action the loop should take for a lead at the given reminder count. */
enum Action { INJECT, WAIT_BUSY, STOP }
/**
* 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
* 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.
*/
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;
}
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);
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);
}
return result;
}
/** 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) {
}
/**
* 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);
}
/**
* 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.
*/
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) {
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());
}
/**
* 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.
* 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
*/
Action decideTickets(String lead, int reminderCount) {
if (pendingFor(lead).isEmpty()) {
log.debug("push: nothing pending for lead {}, stopping ticket reminder", lead);
Action decide(String lead, int replyReminderCount, int ticketReminderCount) {
boolean hasReplyWork = !pendingReplyTargetsFor(lead).isEmpty();
boolean hasTicketWork = !pendingTicketIdsFor(lead).isEmpty();
if (!hasReplyWork && !hasTicketWork) {
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
return Action.STOP;
}
if (reminderCount >= maxReminders) {
log.debug("push: ticket reminder cap ({}) reached for lead {}, stopping", maxReminders, lead);
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
boolean ticketEligible = hasTicketWork && ticketReminderCount < maxReminders;
if (!replyEligible && !ticketEligible) {
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
maxReminders, lead);
countNudge("exhausted");
return Action.STOP;
}
@@ -278,95 +189,194 @@ public final class ReplyPushLoop {
return Action.WAIT_BUSY;
}
/** 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 -> {
injectTicketNudge(lead, reminderCount);
scheduleTicketTick(lead, reminderCount + 1);
}
case WAIT_BUSY -> scheduleTicketTick(lead, reminderCount);
case STOP -> stopOrRestartTicketLoop(lead, pendingBefore);
}
}
// --- public entrypoints ----------------------------------------------------------------------
/** 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());
/**
* 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());
}
/**
* 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).
* 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>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>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);
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);
}
// 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);
}
}
/**
* 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).
*
* <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>Package-private so a test can drive the interleaving directly rather than trying to force a
* 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.
* genuine thread race.
*
* <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.
* <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.
*/
void stopOrRestartTicketLoop(String lead, Set<String> pendingBefore) {
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore) {
activeLeads.remove(lead);
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);
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);
return;
}
log.debug("push: ticket reminder loop ended for lead {}", lead);
log.debug("push: reminder loop ended for lead {}", 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);
/** 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);
return;
}
String nudge = formatTicketsNudge(pending);
String nudge = formatNudge(replyTargets, tickets);
try {
agents.send(lead, nudge);
log.debug("push: ticket nudge {}/{} sent to lead {} for {} ticket(s)",
reminderCount + 1, maxReminders, lead, pending.size());
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}; {} reply target(s), {} ticket(s))",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
replyTargets.size(), tickets.size());
countNudge("delivered");
} catch (RuntimeException e) {
log.warn("push: failed to nudge lead {} for {} ticket(s) (reminder {}/{}): {}",
lead, pending.size(), reminderCount + 1, maxReminders, e.toString());
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}): {}",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, e.toString());
}
}
/** Schedule the next ticket-loop tick on the scheduler thread pool. */
private void scheduleTicketTick(String lead, int nextReminderCount) {
scheduler.schedule(() -> ticketTick(lead, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS);
/** 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);
}
/** Render one or several pending tickets as a single nudge line. */
// --- 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. */
private static String formatTicketsNudge(List<PendingTicket> pending) {
if (pending.size() == 1) {
PendingTicket t = pending.get(0);
@@ -383,23 +393,22 @@ public final class ReplyPushLoop {
// --- lifecycle -----------------------------------------------------------------------------
/**
* 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.
* 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.
*/
public boolean isActive() {
return !activeTargets.isEmpty() || !activeLeads.isEmpty();
return !activeLeads.isEmpty();
}
/** Shut down the scheduler. Outstanding reminders are cancelled. */
public void stop() {
scheduler.shutdownNow();
activeTargets.clear();
activeLeads.clear();
pendingReplies.clear();
pendingTickets.clear();
}
@@ -13,6 +13,7 @@ import dev.ltms.bridged.herdr.HerdrClient;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.placement.PlacementException;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.peer.MemberRole;
@@ -292,6 +293,11 @@ public final class BridgedApp {
ctx.status(201).json(view(member));
} catch (GuardException e) {
ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage()));
} catch (PlacementException e) {
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — a benign,
// likely-transient refusal, distinct from "profile does not exist" below. 503: the
// request was valid and will likely succeed later.
ctx.status(503).json(Map.of("error", "no_capacity", "detail", e.getMessage()));
} catch (IllegalArgumentException e) {
ctx.status(400).json(Map.of("error", "unknown_profile", "detail", e.getMessage()));
} catch (PeerUnreachableException e) {
@@ -16,12 +16,15 @@ import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.session.WorktreeRequest;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.member.CompositePeerLauncher;
import dev.ltms.bridged.placement.BackendQuarantine;
import dev.ltms.bridged.placement.PlacementPolicies;
import io.modelcontextprotocol.spec.McpSchema;
import dev.ltms.bridged.msg.InMemoryReplyInbox;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
@@ -423,6 +426,35 @@ class BridgeMcpTest {
assertTrue(textOf(res).contains("unknown worker profile"), textOf(res));
}
/**
* CB-599: a profile at its {@code maxLoad} cap must surface a readable reason on the MCP
* surface too, not merely flip {@code isError} with an opaque or absent message.
*/
@Test
void spawnAtMaxLoadSurfacesTheCapacityReason() {
FakeHerdr h = new FakeHerdr();
BridgedConfig.Profile wcfg = new BridgedConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN", null,
"tab", "bridged-workers", "worker: {profile} #{n}", null,
null, null, null, null, null, null, null, 0, null, null, null);
Map<String, BridgedConfig.Profile> profiles = Map.of(wcfg.profile(), wcfg);
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(
new AgentControl(h), new WorkspaceControl(h), new SubscriptionGuard(Set.of("gx00.gw")),
profiles, wcfg.profile(), k -> "BRIDGED_WORKER_TOKEN".equals(k) ? "tok" : null);
CompositePeerLauncher composite = new CompositePeerLauncher(
List.of(delegate), wcfg.profile(), profiles, PlacementPolicies.fixed(), _ -> 0);
SessionManager sm = new SessionManager(composite);
McpSchema.CallToolResult res = BridgeMcp.spawn(sm, "ltms-local");
assertTrue(res.isError());
String text = textOf(res);
assertTrue(text.contains("no capacity"), "surfaces a capacity reason, not a bare error: " + text);
assertTrue(text.contains("ltms-local"), "names the profile: " + text);
assertTrue(text.contains("maxLoad"), "explains the refusal: " + text);
assertFalse(h.called("agent.start"), "at cap, the spawn is refused before any herdr call");
}
@Test
void spawnPassesTheRequestedCwdToTheWorker() {
FakeHerdr h = new FakeHerdr();
@@ -16,6 +16,7 @@ import java.util.List;
import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
@@ -166,6 +167,88 @@ class AmqpReplyInboxContractTest {
}
}
@Test
void prefetchBoundsTheHeldBacklog() throws Exception {
String target = "worker-prefetch-" + System.nanoTime();
int prefetch = 4;
int published = 10;
try (AmqpReplyInbox inbox = new AmqpReplyInbox(newConnection(), prefetch);
Connection inspect = newConnection()) {
inbox.own(target);
for (int i = 0; i < published; i++) {
inbox.publish(target, "m" + i, "payload " + i);
}
awaitHeldAtLeast(inbox, target, prefetch);
long depth;
try (Channel ch = inspect.createChannel()) {
depth = ch.queueDeclarePassive(queueName(target)).getMessageCount();
}
assertTrue(depth >= published - prefetch,
"broker should still hold at least " + (published - prefetch)
+ " undelivered messages behind a prefetch of " + prefetch + ", saw " + depth);
drainUntilEmpty(inbox, target, inspect);
}
}
@Test
void unroutablePublishReportsFailureNotSilentSuccess() throws Exception {
String target = "worker-unroutable-" + System.nanoTime();
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
// Deliberately never own(target): the queue is never declared, so the default-exchange
// route to agent.<target>.inbox does not exist and the broker must return the publish.
IllegalStateException ex = assertThrows(IllegalStateException.class,
() -> inbox.publish(target, "m1", "nobody home"));
assertTrue(ex.getMessage() != null && ex.getMessage().toLowerCase().contains("unroutable"),
"expected an unroutable-publish failure, got: " + ex.getMessage());
}
}
@Test
void confirmedPublishDeliversNormally() throws Exception {
String target = "worker-confirm-" + System.nanoTime();
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
inbox.own(target);
inbox.publish(target, "m1", "confirmed delivery"); // must return normally: routed and confirmed
List<ReplyInbox.InboxMessage> got = awaitPeek(inbox, target);
assertEquals(1, got.size());
assertEquals("confirmed delivery", got.getFirst().content());
inbox.ack(target, "m1");
}
}
/** Poll peek until at least {@code n} replies for {@code target} are held, or ~10s elapse. */
@SuppressWarnings("BusyWait")
private static void awaitHeldAtLeast(AmqpReplyInbox inbox, String target, int n) throws InterruptedException {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
while (inbox.peek(target).size() < n && System.nanoTime() < deadline) {
Thread.sleep(50);
}
}
/**
* Repeatedly ack whatever is currently held (freeing prefetch slots for the next delivery) until
* both the local snapshot and the broker's own queue depth are empty, or ~10s elapse.
*/
@SuppressWarnings("BusyWait")
private static void drainUntilEmpty(AmqpReplyInbox inbox, String target, Connection inspect) throws Exception {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
while (System.nanoTime() < deadline) {
for (ReplyInbox.InboxMessage msg : inbox.peek(target)) {
inbox.ack(target, msg.msgId());
}
try (Channel ch = inspect.createChannel()) {
if (ch.queueDeclarePassive(queueName(target)).getMessageCount() == 0 && inbox.peek(target).isEmpty()) {
return;
}
}
Thread.sleep(50);
}
throw new AssertionError("did not drain " + target + " to empty within the deadline");
}
/** Poll peek (broker delivery is async) until a reply for {@code target} appears or ~10s elapse. */
@SuppressWarnings("BusyWait") // deliberate poll for async broker delivery, bounded by the deadline
private static List<ReplyInbox.InboxMessage> awaitPeek(AmqpReplyInbox inbox, String target)
@@ -0,0 +1,264 @@
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,6 +20,7 @@ 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.*;
@@ -27,6 +28,12 @@ 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.
@@ -58,38 +65,50 @@ class ReplyPushLoopTest {
// --- decide() logic ------------------------------------------------------------------------
@Test
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));
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");
}
@Test
void decideWithEmptyInboxIsStop() {
void decideWithNothingPendingIsStop() {
agents = agentWithStatus("idle");
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0));
}
@Test
void decideAtCapIsStop() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100).decide(WORKER, 2));
var loop = loop(2, 100);
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0));
}
@Test
void decideUnderCapWithInjectablePrimaryIsInject() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
}
@Test
void decideUnderCapWithBlockedPrimaryIsInject() {
agents = agentWithStatus("blocked");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
"BLOCKED is injectable");
}
@@ -97,7 +116,9 @@ class ReplyPushLoopTest {
void decideUnderCapWithDonePrimaryIsInject() {
agents = agentWithStatus("done");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
"DONE is injectable");
}
@@ -105,23 +126,29 @@ class ReplyPushLoopTest {
void decideUnderCapWithBusyPrimaryIsWaitBusy() {
agents = agentWithStatus("working");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
}
@Test
void decideUnderCapWithUnknownPrimaryIsWaitBusy() {
agents = agentWithStatus("unknown");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
}
@Test
void decideStopsAfterInboxIsEmptied() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
inbox.ack(WORKER, "m1");
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0));
}
// --- onReplyQueued integration -------------------------------------------------------------
@@ -201,12 +228,19 @@ class ReplyPushLoopTest {
assertTrue(nudge.contains("bridge_poll(target=term_worker)"));
}
// --- CB-588: async ticket terminal nudges — decideTickets() logic --------------------------
@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 -----------------------
@Test
void decideTicketsWithNothingPendingIsStop() {
agents = agentWithStatus("idle");
assertEquals(ReplyPushLoop.Action.STOP, loop().decideTickets(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0));
}
@Test
@@ -214,7 +248,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
var loop = loop(2, 100_000);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 2));
}
@Test
@@ -222,7 +256,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
var loop = loop(5, 100_000);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decideTickets(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
}
@Test
@@ -230,7 +264,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("working");
var loop = loop(5, 100_000);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decideTickets(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
}
@Test
@@ -238,7 +272,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 100_000);
loop.onTicketTerminal("task-1", WORKER, false); // no lead known -> never registered as pending
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0));
}
// --- CB-588: onTicketTerminal integration ---------------------------------------------------
@@ -255,7 +289,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 nudge must not tell the lead to run the target-poll call");
assertFalse(nudge.contains("bridge_poll(target="), "a ticket-only nudge must not tell the lead to run the target-poll call");
}
@Test
@@ -346,11 +380,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
// decideTickets' STOP releases it, so the ticket coalesces onto a schedule that is about to
// decide's 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 stopOrRestartTicketLoop — the STOP path's own release-and-recheck —
// reliable, so this drives stopOrRestart — 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 (pendingBefore=empty), standing in for one
// that was NOT part of the pre-decision snapshot (ticketsBefore=empty), standing in for one
// that races in during the decision-to-release window.
var rec = recordingClient();
agents = new AgentControl(rec);
@@ -358,9 +392,9 @@ class ReplyPushLoopTest {
loop.onTicketTerminal("task-1", WORKER, false); // pendingTickets={task-1}; activeLeads={PRIMARY}
// Stand in for the scheduler thread reaching decideTickets==STOP for this lead — with nothing
// Stand in for the scheduler thread reaching decide==STOP for this lead — with nothing
// pending at decide time — while task-1 races in before the release below runs.
loop.stopOrRestartTicketLoop(PRIMARY, Set.of());
loop.stopOrRestart(PRIMARY, Set.of(), 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");
@@ -368,23 +402,67 @@ class ReplyPushLoopTest {
@Test
void aStaleUncollectedTicketAtCapDoesNotRestartTheLoop() {
// The other direction of the same fix: restarting on ANY non-empty pendingFor(lead) would be
// The other direction of the same fix: restarting on ANY non-empty pending set 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
// #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.
// #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.
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.stopOrRestartTicketLoop(PRIMARY, Set.of("task-1"));
loop.stopOrRestart(PRIMARY, Set.of(), 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");
@@ -402,7 +480,103 @@ class ReplyPushLoopTest {
void inboxNudgeStillUsesTheOriginalTargetPollCall() {
String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER);
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
"CB-588 must not change the CB-307 inbox nudge's call shape");
"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");
}
// --- metrics (CB-512) ----------------------------------------------------------------------
@@ -430,8 +604,10 @@ 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(2, 100, metrics).decide(WORKER, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0));
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the reminder cap must count as exhausted");
@@ -459,7 +635,7 @@ class ReplyPushLoopTest {
var loop = loop(2, 100_000, metrics);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 2));
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the ticket reminder cap must count as exhausted");
@@ -547,4 +723,49 @@ 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() {
}
}
}
@@ -18,6 +18,8 @@ import dev.ltms.bridged.session.GitWorktrees;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.session.Worktrees;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.member.CompositePeerLauncher;
import dev.ltms.bridged.placement.PlacementPolicies;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
@@ -240,6 +242,43 @@ class BridgedAppTest {
assertFalse(herdr.called("agent.start"), "an unknown profile must not spawn anything");
}
/**
* CB-599: a profile at its {@code maxLoad} cap must not surface as a bare 500 — the caller
* needs a structured, readable reason, distinct from "unknown_profile".
*/
@Test
void spawnAtMaxLoadIs503WithTheCapacityReasonNotABare500() throws Exception {
FakeHerdr herdr = new FakeHerdr();
BridgedConfig.Profile wcfg = new BridgedConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN", null,
"tab", "bridged-workers", "worker: {profile} #{n}", null,
null, null, null, null, null, null, null, 0, null, null, null);
Map<String, BridgedConfig.Profile> profiles = Map.of(wcfg.profile(), wcfg);
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(
new AgentControl(herdr), new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
profiles, wcfg.profile(), k -> "BRIDGED_WORKER_TOKEN".equals(k) ? "tok-abc" : null);
CompositePeerLauncher workers = new CompositePeerLauncher(
List.of(delegate), wcfg.profile(), profiles, PlacementPolicies.fixed(), _ -> 0);
SessionManager sessions = new SessionManager(workers, new GitWorktrees());
this.presence = sessions.asPresence();
Injector injector = new Injector(new AgentControl(herdr));
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
sessions.onAcquire(inbox::own);
MessageService messages = new MessageService(new AgentControl(herdr), injector, new Rendezvous(), inbox);
app = new BridgedApp(herdr, workers, sessions, messages, this.presence, null).build().start("127.0.0.1", 0);
int port = app.port();
HttpResponse<String> res = req(port, "POST", "/members?profile=ltms-local");
assertEquals(503, res.statusCode(), res.body());
JsonNode body = mapper.readTree(res.body());
assertEquals("no_capacity", body.get("error").asText());
String detail = body.get("detail").asText();
assertTrue(detail.contains("ltms-local"), "detail names the profile: " + detail);
assertTrue(detail.contains("maxLoad"), "detail explains the refusal: " + detail);
assertFalse(herdr.called("agent.start"), "at cap, the spawn is refused before any herdr call");
}
@Test
void spawnWorkerReusesExistingWorkerSpace() throws Exception {
// A space labelled "bridged-workers" already exists → no second workspace.create.