Compare commits

..

5 Commits

Author SHA1 Message Date
Dai Ha 32bf324a1e CB-602: guard against a config key that never reaches the example
CI / contract (pull_request) Successful in 49s
CI / build (pull_request) Failing after 1m24s
BridgedConfig.KNOWN_TOP_LEVEL_KEYS is now package-private so a test can assert
every key the parser accepts appears in bridged.example.yaml — live or
commented-out, since the file is gitignored and the example is the only
committed description of the config schema. The existing tests only checked
the example->code direction; this adds code->example.
2026-08-16 18:06:36 +02:00
ltms 65c9deb4d1 Merge CB-597: correct bridged.example.yaml, including two knobs that do nothing
CI / contract (push) Successful in 42s
CI / build (push) Successful in 2m2s
My premise for this ticket was wrong and the worker corrected it. I reported five whole sections missing from the example; nothing was missing. My comparison script only counted uncommented lines, so every section documented as a commented-out example looked absent. Earlier tickets had each updated the example alongside their feature.

What it found instead is more useful than what I asked for — real inaccuracies, found by tracing each field through the parser and its consumers:

- `health.workingSuspectAfterSeconds` and `paneProbeIntervalSeconds` documented enforced minimums that do not exist. I checked: both names appear **only** in the `Health` record declaration and are read by nothing. Only `intervalSeconds` is clamped, and it is silently raised to 15 rather than rejected.
- `notifications.mode: webhook` only flips what `bridge_list` reports as `healthCoverage`. It sends no webhook — "webhook" appears in one `configured()` boolean and there is no delivery code in the repo.
- `lifecycle.clearAfterTurn` was undocumented, and is a no-op for any peer kind other than claude-code.
- The reload doc claimed the whole `fleet:` block is hot; `fleet.leaders` is built once at startup and is not rebuilt, so a change is silently accepted and does nothing until a restart.
- The `fleet.leaders` demotion consequence is now stated next to the block itself: an unmatched pane is silently an ordinary worker and every orchestration call it makes is refused, with no startup error.

Documenting a knob as dead is worth more than documenting it as working. Someone tuning `workingSuspectAfterSeconds` would otherwise have concluded their monitor was broken.

Comments only — no parsing or production code touched. Verified by the lead: parses cleanly under the project's own snakeyaml 1.30, top-level live keys `[bind, herdrSocket, profiles, placement, fleet, guard]`, the rest correctly commented examples. Both dead-knob claims verified by grep against `src/main` rather than taken on the worker's word.
2026-08-16 17:59:41 +02:00
ltms 28ae27b8e1 Merge CB-599: a capacity refusal now tells the caller why
CI / build (push) Failing after 1m19s
CI / contract (push) Successful in 1m25s
`PlacementException extends IllegalStateException`, and neither spawn path caught that type, so it escaped to Javalin's default handler as a bare `500 Server Error` with a text/plain body — while every other failure on the same endpoint returned structured JSON. The reason existed and was good, but only in the daemon log.

I hit this live while orchestrating: asked for a member on a full profile, got a blank 500, guessed another profile, got a blank 500 again, and spent two round trips learning things the daemon already knew.

Both surfaces now catch it. REST returns 503 with `{"error":"no_capacity","detail":...}`; MCP returns the same reason in the `isError` shape it already uses for every other spawn failure. 503 is right because the request was valid and will likely succeed later — the caller did nothing wrong, so 400 would have been a lie.

The other throw sites all funnel through the same type, so quarantine cooldowns, weight-0 exclusion, and the all-at-cap / all-quarantined / all-unreachable messages now reach callers too. That last group matters most: those three distinguish "wait a moment" from "your backends are gone", and all three used to arrive as the identical blank 500.

Tests assert the caller can read the *reason*, not merely that the status changed — one per surface.

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

No exception message was reworded. This change delivers messages that were already written.
2026-08-16 17:54:44 +02:00
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
Dai Ha 0a2b3a4a56 CB-597: fix inaccuracies the example config already had, none actually missing
CI / contract (pull_request) Successful in 1m7s
CI / build (pull_request) Failing after 1m21s
Audited bridged.example.yaml against BridgedConfig's KNOWN_TOP_LEVEL_KEYS and
found every top-level key already documented (broker, health, lifecycle,
configReload, quarantineCooldownSeconds, fleet.leaders/architects/reviewers
included) — CB-573/CB-566/CB-559/CB-579/CB-527/528 each updated the example
alongside their feature. What was actually wrong:

- health.workingSuspectAfterSeconds/paneProbeIntervalSeconds claimed enforced
  minimums (300/60) that don't exist in code — only intervalSeconds is
  clamped (floor 15); the other two are parsed but never read anywhere.
- notifications.mode: webhook was undocumented as only flipping the
  healthCoverage label bridge_list reports — no webhook is ever sent.
- lifecycle.clearAfterTurn was missing entirely.
- the HOT bullet under configReload claimed the whole fleet: block reloads
  live, but ConfigRef's own javadoc carves out fleet.leaders as needing a
  restart with no deferred-list warning — added that exception.
- fleet.leaders' demotion consequence (unmatched tab -> silent WORKER
  demotion, no startup error) is now stated inline next to the block, not
  just implied by the multi-lead rationale higher up.
2026-08-16 17:37:28 +02:00
9 changed files with 222 additions and 191 deletions
+35 -6
View File
@@ -88,15 +88,28 @@ bind:
# backoffMs: 60000
# quietNudgeCap: 3
# Fleet health detection is dormant unless enabled. It reads one whole-fleet agent list per tick.
# It can run without a webhook; bridge_list then reports healthCoverage: detection-only.
# Fleet health detection is dormant unless enabled (CB-573). It reads one whole-fleet agent list
# per tick.
# intervalSeconds → how often a tick runs (default 30). ENFORCED floor of 15: the code computes
# Math.max(15, intervalSeconds), so a lower value is silently raised, not
# rejected.
# workingSuspectAfterSeconds, paneProbeIntervalSeconds → accepted and parsed, but NOT YET READ by
# anything — the dormant monitor only consumes intervalSeconds today (CB-573
# shipped ahead of the evidence publishers these two knobs are for). Setting
# them changes nothing right now, and no minimum is enforced on either, because
# nothing reads them to enforce one. They exist so a later build can start
# honouring them without another config-shape change.
# notifications.mode → "webhook" flips what bridge_list REPORTS (healthCoverage: "full" instead
# of "detection-only") — it does NOT make bridged send any webhook call; no
# delivery mechanism is implemented yet. Any other value, or omitting the
# block, reports "detection-only".
# health:
# enabled: true
# intervalSeconds: 30 # minimum 15
# workingSuspectAfterSeconds: 600 # minimum 300
# paneProbeIntervalSeconds: 60 # minimum 60
# intervalSeconds: 30
# workingSuspectAfterSeconds: 600
# paneProbeIntervalSeconds: 60
# notifications:
# mode: disabled # disabled (default) or webhook
# mode: disabled
# herdr Unix socket. Omit to use the client default
# (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}).
@@ -286,6 +299,11 @@ placement: weighted
# / credentialId. Those are hot because the placement policy (and, for credentialId,
# the CB-578 stage B quarantine check) reads them through a supplier — being config is
# not by itself enough to make a key hot.
# EXCEPT `fleet.leaders`: Bridged.main reads it once at startup to build the lead tab
# scanner and launcher, and neither is rebuilt on reload. A changed/added/removed
# `fleet.leaders` entry is silently accepted — the reload reports "config reloaded"
# with nothing in the deferred list — but has NO effect until you restart. Treat it
# as deferred in practice, even though today's reload output does not say so.
# DEFERRED → accepted into the new config, but the wiring built at startup keeps the old value
# until you restart: `lifecycle:`, `leadHeartbeat:`, `guard:`, `worktreeRoot:`,
# `spawnReadyTimeoutMs` / `spawnReadyPollMs`, `quarantineCooldownSeconds` (CB-578
@@ -368,6 +386,12 @@ fleet:
# An auto-launched lead is NOT a member: it gets no worker reply charter, is never registered with
# the session lifecycle (the idle reaper would kill your orchestrator), and stays on the
# subscription — ANTHROPIC_BASE_URL/AUTH_TOKEN are stripped from its env whatever the profile says.
#
# GET THE `tab:` VALUE RIGHT. A pane that does not match any configured `tab:` (a typo, a renamed
# tab, a pane no entry names at all) is not recognised as a lead — it resolves as an ordinary
# WORKER instead, silently, and every orchestration call it makes (spawn/stop/send/drain) is
# refused. There is no error at startup for this: an unmatched pane is simply not a lead. If your
# primary suddenly can't spawn or send, check this section first.
# leaders:
# opus-5.0:
# profile: opus # omit to never create this lead, only recognise it
@@ -423,10 +447,15 @@ guard:
# idleTtlSeconds → reap READY/DONE sessions idle longer than this (never BUSY/SPAWNING)
# contextCap → force-release a session after this many delegated turns
# drainTimeoutSeconds → seconds to wait for BUSY sessions on shutdown before forced teardown
# clearAfterTurn → whether a reusable worker discards its conversation context after every
# completed delegated turn (default false). Works for claude-code workers
# only — any other peer kind (e.g. opencode) logs "context reset is
# unsupported for peer kind …" once and the reset is a no-op.
# lifecycle:
# idleTtlSeconds: 300
# contextCap: 10
# drainTimeoutSeconds: 5
# clearAfterTurn: false
# Durable reply delivery (CB-307 Stage 2). OMIT this block entirely to keep the default
# in-memory, soft-state reply inbox (late worker replies are held only until a daemon bounce).
@@ -967,8 +967,12 @@ public record BridgedConfig(
/**
* Top-level keys this version understands. Used only to warn about the rest — see
* {@link #warnUnknownTopLevelKeys}. Keep in step with the record components.
*
* <p>Package-private (not {@code private}) so a test can assert every key here is documented in
* {@code bridged.example.yaml} — the only committed description of the config schema, since
* {@code bridged.yaml} itself is gitignored.
*/
private static final Set<String> KNOWN_TOP_LEVEL_KEYS = Set.of(
static final Set<String> KNOWN_TOP_LEVEL_KEYS = Set.of(
"bind", "herdrSocket", "profiles", "guard", "worktreeRoot",
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds");
@@ -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) {
@@ -66,13 +66,8 @@ public final class ReplyPushLoop {
private final long backoffMs;
private final Metrics metrics; // CB-512: nullable — no registry in unit tests
/**
* Worker targets with a reply queued, keyed by target. Each entry carries its own nudge
* count (CB-598) rather than sharing one counter per lead per source: a target's count only
* ever reflects nudges that actually named that target, so a target that joins while the
* schedule is already deep into another target's reminders still reads as fresh.
*/
private final ConcurrentHashMap<String, ReplyEntry> pendingReplies = new ConcurrentHashMap<>();
/** 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-590: leads with an active combined reminder schedule (replies and/or tickets). */
@@ -121,10 +116,10 @@ public final class ReplyPushLoop {
Set<String> result = new HashSet<>();
for (var entry : pendingReplies.entrySet()) {
String target = entry.getKey();
ReplyEntry owning = entry.getValue();
if (!lead.equals(owning.lead())) continue;
String owningLead = entry.getValue();
if (!lead.equals(owningLead)) continue;
if (inbox.peek(target).isEmpty()) {
pendingReplies.remove(target, owning);
pendingReplies.remove(target, owningLead);
continue;
}
result.add(target);
@@ -132,15 +127,8 @@ public final class ReplyPushLoop {
return result;
}
/** A pending reply target: which lead to nudge, and how many nudges have named it so far. */
private record ReplyEntry(String lead, int nudgeCount) {
}
/**
* A ticket awaiting collection: which lead to nudge, whether it ended in failure, and how
* many nudges have named it so far (CB-598 — tracked per ticket, not per lead per source).
*/
private record PendingTicket(String ticket, String lead, boolean failed, int nudgeCount) {
/** A ticket awaiting collection: which lead to nudge, and whether it ended in failure. */
private record PendingTicket(String ticket, String lead, boolean failed) {
}
/** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */
@@ -154,41 +142,6 @@ public final class ReplyPushLoop {
.collect(Collectors.toUnmodifiableSet());
}
/**
* The reply-source reminder count {@link #decide} should see for {@code lead} on this tick:
* the <em>minimum</em> nudge count among the reply targets currently pending for it (CB-598).
*
* <p>Before this, the count passed to {@code decide} was a single counter carried forward
* across scheduled ticks ({@code scheduleNext(lead, count + 1, ...)}), incremented whenever
* the source had <em>any</em> pending work — not tied to which target that work was. A target
* that joined while an older target's count was already near the cap inherited that count on
* its very next tick, even though no nudge had ever named it. Taking the minimum over what is
* actually pending now means a fresh target (count 0) keeps the source eligible regardless of
* how many times an older, still-undrained target has already been nudged; that older target
* keeps riding along in the combined nudge text without spending any more of its own budget
* (see {@link #bumpNudgeCounts}). Returns 0 when nothing is pending — {@link #decide} never
* consults the count in that case, since {@code hasReplyWork} is false.
*/
private int minReplyNudgeCountFor(String lead) {
int min = Integer.MAX_VALUE;
for (String target : pendingReplyTargetsFor(lead)) {
ReplyEntry entry = pendingReplies.get(target);
if (entry != null) {
min = Math.min(min, entry.nudgeCount());
}
}
return min == Integer.MAX_VALUE ? 0 : min;
}
/** As {@link #minReplyNudgeCountFor}, for the ticket source. */
private int minTicketNudgeCountFor(String lead) {
int min = Integer.MAX_VALUE;
for (PendingTicket ticket : pendingTicketsFor(lead)) {
min = Math.min(min, ticket.nudgeCount());
}
return min == Integer.MAX_VALUE ? 0 : min;
}
/**
* Pure decision function: examine everything pending for {@code lead} — reply targets and
* tickets alike — and return what the loop should do.
@@ -202,14 +155,9 @@ public final class ReplyPushLoop {
* {@link Action#INJECT}. Only when neither source has eligible work does the loop
* {@link Action#STOP}.
*
* <p><strong>CB-598: the counts are per-item, not per-tick.</strong> {@link #tick} no longer
* carries these counts forward across scheduled calls — it recomputes them fresh every tick via
* {@link #minReplyNudgeCountFor} / {@link #minTicketNudgeCountFor}, so this function itself did
* not need to change; only what its caller feeds it did.
*
* @param lead the lead terminal to nudge
* @param replyReminderCount the lowest nudge count among reply targets pending for this lead
* @param ticketReminderCount the lowest nudge count among tickets pending for this lead
* @param replyReminderCount how many nudges have covered pending reply work for this lead
* @param ticketReminderCount how many nudges have covered pending ticket work for this lead
* @return the action the caller should take
*/
Action decide(String lead, int replyReminderCount, int ticketReminderCount) {
@@ -256,8 +204,7 @@ public final class ReplyPushLoop {
log.debug("push: no lead is known to be waiting on {}, skipping reminder", target);
return;
}
pendingReplies.compute(target, (t, existing) ->
new ReplyEntry(lead.get(), existing == null ? 0 : existing.nudgeCount()));
pendingReplies.put(target, lead.get());
startOrCoalesce(lead.get());
}
@@ -285,8 +232,7 @@ public final class ReplyPushLoop {
ticket, target);
return;
}
pendingTickets.compute(ticket, (id, existing) ->
new PendingTicket(ticket, lead.get(), failed, existing == null ? 0 : existing.nudgeCount()));
pendingTickets.put(ticket, new PendingTicket(ticket, lead.get(), failed));
startOrCoalesce(lead.get());
}
@@ -309,37 +255,28 @@ public final class ReplyPushLoop {
return;
}
log.debug("push: starting reminder loop for lead {}", lead);
scheduleNext(lead);
scheduleNext(lead, 0, 0);
}
/**
* Execute one loop tick — called on the scheduler thread (or directly by a test; package-private
* for the same reason as {@link #stopOrRestart}).
*
* <p><strong>CB-598.</strong> The reminder counts fed into {@link #decide} are recomputed fresh
* every tick from what is actually pending right now ({@link #minReplyNudgeCountFor} /
* {@link #minTicketNudgeCountFor}), rather than carried forward as running counters across
* scheduled calls. A counter carried forward has no memory of which item it was counting for:
* a target or ticket that joined mid-backoff — after the previous tick fired but before this one
* did — is already sitting in {@code repliesBefore} / {@code ticketsBefore} below by the time this
* tick takes its snapshot, indistinguishable at that point from backlog the cap is meant to
* silence. Recomputing from the per-item counts fixes that: a newly-joined item's own count is
* still 0, so it keeps its source eligible regardless of how depleted an older, still-undrained
* item's count is.
*/
void tick(String lead) {
/** 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);
int replyReminderCount = minReplyNudgeCountFor(lead);
int ticketReminderCount = minTicketNudgeCountFor(lead);
var action = decide(lead, replyReminderCount, ticketReminderCount);
switch (action) {
case INJECT -> {
injectNudge(lead, replyReminderCount, ticketReminderCount);
scheduleNext(lead);
// 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);
case WAIT_BUSY -> scheduleNext(lead, replyReminderCount, ticketReminderCount);
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore);
}
}
@@ -380,7 +317,7 @@ public final class ReplyPushLoop {
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t));
if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
scheduleNext(lead);
scheduleNext(lead, 0, 0);
return;
}
log.debug("push: reminder loop ended for lead {}", lead);
@@ -407,28 +344,11 @@ public final class ReplyPushLoop {
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}): {}",
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, e.toString());
}
// Bump every item actually named in this nudge, not just whatever the shared source-level
// eligibility used to gate (CB-598) — each item's own count is what the next tick's
// minReplyNudgeCountFor / minTicketNudgeCountFor will read. An item already at or over the
// cap keeps riding along in the text (still pending, still named) but its extra bumps here
// are inert: decide() already treats it as ineligible once its count reaches maxReminders.
bumpNudgeCounts(replyTargets, tickets);
}
/** Record that every one of these items was just named in a sent (or attempted) nudge. */
private void bumpNudgeCounts(Set<String> replyTargets, List<PendingTicket> tickets) {
for (String target : replyTargets) {
pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1));
}
for (PendingTicket ticket : tickets) {
pendingTickets.computeIfPresent(ticket.ticket(),
(id, e) -> new PendingTicket(e.ticket(), e.lead(), e.failed(), e.nudgeCount() + 1));
}
}
/** Schedule the next tick on the scheduler thread pool. */
private void scheduleNext(String lead) {
scheduler.schedule(() -> tick(lead),
private void scheduleNext(String lead, int nextReplyReminderCount, int nextTicketReminderCount) {
scheduler.schedule(() -> tick(lead, nextReplyReminderCount, nextTicketReminderCount),
backoffMs, TimeUnit.MILLISECONDS);
}
@@ -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) {
@@ -10,6 +10,7 @@ import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.regex.Pattern;
import static org.junit.jupiter.api.Assertions.*;
@@ -1189,6 +1190,79 @@ class BridgedConfigTest {
assertEquals(5, cfg.leadHeartbeat().quietNudgeCap());
}
/**
* A top-level key {@code BridgedConfig} reads but that appears nowhere in
* {@code bridged.example.yaml} — live or commented — is invisible drift: {@code bridged.yaml}
* is gitignored, so the example is the ONLY committed description of the config schema, and
* neither {@link #shippedExampleConfigParses} (example → code: does the example still parse)
* nor {@link #everyOptionalKnobDocumentedInTheExampleBinds} (a hand-maintained list of keys
* that must bind) can catch a brand-new key nobody added to either.
*
* <p>This test compares the OTHER direction: every key in {@link BridgedConfig#KNOWN_TOP_LEVEL_KEYS}
* (the parser's own accepted set, which backs the unknown-key WARN) must appear as a top-level
* key in the example text, live or commented-out — see {@link #topLevelKeyDocumented}.
*/
@Test
void everyKnownTopLevelKeyIsDocumentedInTheExample() throws Exception {
Path example = Path.of("bridged.example.yaml");
assertTrue(Files.exists(example), "bridged.example.yaml must ship next to the pom");
String text = Files.readString(example);
List<String> undocumented = BridgedConfig.KNOWN_TOP_LEVEL_KEYS.stream()
.filter(key -> !topLevelKeyDocumented(text, key))
.sorted()
.toList();
assertTrue(undocumented.isEmpty(), () -> "key(s) " + undocumented
+ " are read by BridgedConfig but appear nowhere in bridged.example.yaml — "
+ "document each one there, commented out if optional. bridged.yaml is "
+ "gitignored, so this file is the only committed description of the config "
+ "schema an operator or a worker can see.");
}
/**
* Most of {@code bridged.example.yaml} is deliberately commented out — optional sections are
* documented as commented blocks so the shipped file stays a working minimal config. A key
* documented ONLY as a comment must still count as documented; parsing the file as YAML and
* reading its live key set (as an earlier attempt at this guard did) gets this wrong, because
* every commented section then looks entirely absent.
*/
@Test
void commentedOnlyTopLevelKeyCountsAsDocumented() {
String yaml = """
bind:
port: 8765
# broker:
# uri: amqp://guest:guest@127.0.0.1:5672
""";
assertTrue(topLevelKeyDocumented(yaml, "broker"),
"a key documented only inside a commented-out block must still count as documented");
}
/** A key that appears in neither a live nor a commented top-level line must NOT count. */
@Test
void absentTopLevelKeyIsNotDocumented() {
String yaml = """
bind:
port: 8765
""";
assertFalse(topLevelKeyDocumented(yaml, "broker"),
"a key mentioned nowhere in the example must not be reported as documented");
}
/**
* True when {@code key} appears as a top-level YAML key in {@code yaml} — either live
* ({@code key:} at column 0) or commented out ({@code # key:}, also at column 0, with only
* whitespace between the {@code #} and the key). Anchoring on column 0 is what keeps this a
* top-level check: an indented occurrence (a nested field, or prose inside a comment that
* happens to end in a colon) never matches, because {@code ^} requires the key's own first
* character — or the sole leading {@code #} — to sit at the very start of the line.
*/
private static boolean topLevelKeyDocumented(String yaml, String key) {
Pattern p = Pattern.compile("(?m)^(?:#\\s*)?" + Pattern.quote(key) + ":");
return p.matcher(yaml).find();
}
@Test
void placementDefaultsToFixedForExistingConfigs(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-placement.yaml");
@@ -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();
@@ -42,7 +42,6 @@ class ReplyPushLoopTest {
private static final String PRIMARY = "term_primary";
private static final String WORKER = "term_worker";
private static final String WORKER2 = "term_worker2";
private static final ObjectMapper MAPPER = new ObjectMapper();
private PrimaryRegistry registry;
@@ -580,83 +579,6 @@ class ReplyPushLoopTest {
"both the reply and the ticket source are at their own cap — must still stop");
}
// --- CB-598: work arriving during a backoff must not read as stale backlog -------------------
@Test
void aTargetArrivingDuringTheBackoffGetsNudgedDespiteAnAlreadyCappedSibling() {
// The bug: reminder counts used to be a single counter per lead per source, carried
// forward across scheduled ticks (scheduleNext(lead, count + 1, ...)) rather than tracked
// per pending item. WORKER gets nudged once here, which — with cap=1 — exhausts the
// shared reply-source counter for this lead. WORKER2 then queues a reply for the SAME
// lead "during the backoff": while the schedule from WORKER's tick is still active, before
// the next tick's own start-of-tick snapshot runs. At that next tick, the OLD code passed
// the already-exhausted shared counter into decide() regardless of WORKER2 never having
// been named in any nudge, and — because WORKER2 was already present in that tick's
// "before" snapshot — stopOrRestart's race check (proven correct on its own elsewhere in
// this file) does not save it either: it looks like ordinary stale backlog, not a race.
// WORKER2 was then stranded forever with no live schedule and no nudge ever naming it.
//
// tick() is driven directly (package-private, same reasoning as stopOrRestart being
// directly testable) so the exact interleaving is deterministic instead of racing the
// scheduler thread over a real ~15s backoff.
//
// Before the fix, this test fails on the second assertEquals: rec.sendCount() stays at 1
// (decide() returns STOP on the second tick(), so injectNudge is never called a second
// time) and the "must still get one" assertion never even runs.
int cap = 1;
var rec = recordingClient();
agents = new AgentControl(rec);
inbox.own(WORKER2);
inbox.publish(WORKER, "m1", "hello");
var loop = loop(cap, 100_000); // huge backoff — nothing fires on its own; we drive tick()
loop.onReplyQueued(WORKER);
loop.tick(PRIMARY); // first tick: nudges WORKER alone; WORKER's own count reaches the cap
assertEquals(1, rec.sendCount(), "the first tick should nudge about WORKER");
// WORKER2 "arrives during the backoff": queued for the same lead while the schedule from
// the tick above is still active (activeLeads still holds PRIMARY), before the next tick
// (simulated below) takes its own start-of-tick snapshot.
inbox.publish(WORKER2, "m2", "hello2");
loop.onReplyQueued(WORKER2);
loop.tick(PRIMARY); // the tick that would fire once that backoff elapsed
assertEquals(2, rec.sendCount(),
"WORKER2 was never named in any nudge yet and must still get one, even though "
+ "WORKER's own reminder count is already at the cap");
String secondNudge = rec.sentParams().get(1).getValue().toString();
assertTrue(secondNudge.contains(WORKER2), "the never-named target must be named: " + secondNudge);
// Criterion #3: isActive() must reflect that this lead still had a live nudge to give —
// the second tick took the INJECT branch, so the schedule stayed live rather than being
// torn down under WORKER2.
assertTrue(loop.isActive(), "the schedule must stay active after nudging the fresh target");
}
@Test
void aTicketArrivingDuringTheBackoffGetsNudgedDespiteAnAlreadyCappedSibling() {
// Mirrors the reply-side test above for the ticket source.
int cap = 1;
var rec = recordingClient();
agents = new AgentControl(rec);
var loop = loop(cap, 100_000);
loop.onTicketTerminal("task-1", WORKER, false);
loop.tick(PRIMARY); // first tick: nudges task-1 alone; its count reaches the cap
assertEquals(1, rec.sendCount(), "the first tick should nudge about task-1");
loop.onTicketTerminal("task-2", WORKER, false); // arrives during the backoff, same lead
loop.tick(PRIMARY);
assertEquals(2, rec.sendCount(),
"task-2 was never named in any nudge yet and must still get one, even though "
+ "task-1's reminder count is already at the cap");
String secondNudge = rec.sentParams().get(1).getValue().toString();
assertTrue(secondNudge.contains("task-2"), "the never-named ticket must be named: " + secondNudge);
assertTrue(loop.isActive(), "the schedule must stay active after nudging the fresh ticket");
}
// --- metrics (CB-512) ----------------------------------------------------------------------
@Test
@@ -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.