Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 81a0cf4710 CB-598: track reminder counts per pending item, not per lead per source
CI / build (pull_request) Successful in 1m12s
CI / contract (pull_request) Successful in 1m11s
Work that arrived during the ~15s push_backoff_ms window between two
ticks landed in the pending map before the next tick's start-of-tick
snapshot, so a shared per-lead-per-source counter (carried forward via
scheduleNext(lead, count+1, ...)) already treated it as exhausted
backlog even though no nudge had ever named it. ReplyPushLoop.tick now
recomputes each source's reminder count fresh every tick as the
minimum nudge count among that source's currently pending items, so a
freshly-arrived item (count 0) keeps its source eligible regardless of
how depleted an older, still-undrained sibling's count is. decide()
itself is unchanged.
2026-08-16 17:49:21 +02:00
10 changed files with 194 additions and 232 deletions
+6 -35
View File
@@ -88,28 +88,15 @@ bind:
# backoffMs: 60000
# quietNudgeCap: 3
# Fleet health detection is dormant unless enabled (CB-573). It reads one whole-fleet agent list
# per tick.
# intervalSeconds → how often a tick runs (default 30). ENFORCED floor of 15: the code computes
# Math.max(15, intervalSeconds), so a lower value is silently raised, not
# rejected.
# workingSuspectAfterSeconds, paneProbeIntervalSeconds → accepted and parsed, but NOT YET READ by
# anything — the dormant monitor only consumes intervalSeconds today (CB-573
# shipped ahead of the evidence publishers these two knobs are for). Setting
# them changes nothing right now, and no minimum is enforced on either, because
# nothing reads them to enforce one. They exist so a later build can start
# honouring them without another config-shape change.
# notifications.mode → "webhook" flips what bridge_list REPORTS (healthCoverage: "full" instead
# of "detection-only") — it does NOT make bridged send any webhook call; no
# delivery mechanism is implemented yet. Any other value, or omitting the
# block, reports "detection-only".
# Fleet health detection is dormant unless enabled. It reads one whole-fleet agent list per tick.
# It can run without a webhook; bridge_list then reports healthCoverage: detection-only.
# health:
# enabled: true
# intervalSeconds: 30
# workingSuspectAfterSeconds: 600
# paneProbeIntervalSeconds: 60
# intervalSeconds: 30 # minimum 15
# workingSuspectAfterSeconds: 600 # minimum 300
# paneProbeIntervalSeconds: 60 # minimum 60
# notifications:
# mode: disabled
# mode: disabled # disabled (default) or webhook
# herdr Unix socket. Omit to use the client default
# (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}).
@@ -299,11 +286,6 @@ placement: weighted
# / credentialId. Those are hot because the placement policy (and, for credentialId,
# the CB-578 stage B quarantine check) reads them through a supplier — being config is
# not by itself enough to make a key hot.
# EXCEPT `fleet.leaders`: Bridged.main reads it once at startup to build the lead tab
# scanner and launcher, and neither is rebuilt on reload. A changed/added/removed
# `fleet.leaders` entry is silently accepted — the reload reports "config reloaded"
# with nothing in the deferred list — but has NO effect until you restart. Treat it
# as deferred in practice, even though today's reload output does not say so.
# DEFERRED → accepted into the new config, but the wiring built at startup keeps the old value
# until you restart: `lifecycle:`, `leadHeartbeat:`, `guard:`, `worktreeRoot:`,
# `spawnReadyTimeoutMs` / `spawnReadyPollMs`, `quarantineCooldownSeconds` (CB-578
@@ -386,12 +368,6 @@ fleet:
# An auto-launched lead is NOT a member: it gets no worker reply charter, is never registered with
# the session lifecycle (the idle reaper would kill your orchestrator), and stays on the
# subscription — ANTHROPIC_BASE_URL/AUTH_TOKEN are stripped from its env whatever the profile says.
#
# GET THE `tab:` VALUE RIGHT. A pane that does not match any configured `tab:` (a typo, a renamed
# tab, a pane no entry names at all) is not recognised as a lead — it resolves as an ordinary
# WORKER instead, silently, and every orchestration call it makes (spawn/stop/send/drain) is
# refused. There is no error at startup for this: an unmatched pane is simply not a lead. If your
# primary suddenly can't spawn or send, check this section first.
# leaders:
# opus-5.0:
# profile: opus # omit to never create this lead, only recognise it
@@ -447,15 +423,10 @@ guard:
# idleTtlSeconds → reap READY/DONE sessions idle longer than this (never BUSY/SPAWNING)
# contextCap → force-release a session after this many delegated turns
# drainTimeoutSeconds → seconds to wait for BUSY sessions on shutdown before forced teardown
# clearAfterTurn → whether a reusable worker discards its conversation context after every
# completed delegated turn (default false). Works for claude-code workers
# only — any other peer kind (e.g. opencode) logs "context reset is
# unsupported for peer kind …" once and the reset is a no-op.
# lifecycle:
# idleTtlSeconds: 300
# contextCap: 10
# drainTimeoutSeconds: 5
# clearAfterTurn: false
# Durable reply delivery (CB-307 Stage 2). OMIT this block entirely to keep the default
# in-memory, soft-state reply inbox (late worker replies are held only until a daemon bounce).
@@ -555,10 +555,7 @@ public final class Bridged {
* where set (opt-in). Derived from the config, not hard-coded, so a new profile is covered for
* free. A var required by more than one profile is one entry naming every profile that needs
* it. Deliberately excludes {@code auth.tokenEnv}: that one is already enforced loudly, by a
* startup throw in {@code main()} — about 370 lines <em>below</em> this method's call site
* ({@link #reportRequiredSecrets(BridgedConfig)}), not a few lines above it. That throw only
* fires when {@code auth.mode: token} is configured; under the default loopback-trust mode it
* never runs, and {@code auth.tokenEnv} is simply not required.
* startup throw, a few lines above this method's call site.
*
* <p>Package-private and pure (no I/O, no logging) so the derivation is unit-testable without
* capturing log output; {@link #reportRequiredSecrets(BridgedConfig)} is the logging caller.
@@ -15,7 +15,6 @@ 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;
@@ -693,10 +692,6 @@ 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,8 +66,13 @@ public final class ReplyPushLoop {
private final long backoffMs;
private final Metrics metrics; // CB-512: nullable — no registry in unit tests
/** Worker targets with a reply queued, and the lead to nudge about it, keyed by target. */
private final ConcurrentHashMap<String, String> pendingReplies = new ConcurrentHashMap<>();
/**
* 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<>();
/** 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). */
@@ -116,10 +121,10 @@ public final class ReplyPushLoop {
Set<String> result = new HashSet<>();
for (var entry : pendingReplies.entrySet()) {
String target = entry.getKey();
String owningLead = entry.getValue();
if (!lead.equals(owningLead)) continue;
ReplyEntry owning = entry.getValue();
if (!lead.equals(owning.lead())) continue;
if (inbox.peek(target).isEmpty()) {
pendingReplies.remove(target, owningLead);
pendingReplies.remove(target, owning);
continue;
}
result.add(target);
@@ -127,8 +132,15 @@ public final class ReplyPushLoop {
return result;
}
/** A ticket awaiting collection: which lead to nudge, and whether it ended in failure. */
private record PendingTicket(String ticket, String lead, boolean failed) {
/** 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) {
}
/** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */
@@ -142,6 +154,41 @@ 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.
@@ -155,9 +202,14 @@ 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 how many nudges have covered pending reply work for this lead
* @param ticketReminderCount how many nudges have covered pending ticket work for this lead
* @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
* @return the action the caller should take
*/
Action decide(String lead, int replyReminderCount, int ticketReminderCount) {
@@ -204,7 +256,8 @@ public final class ReplyPushLoop {
log.debug("push: no lead is known to be waiting on {}, skipping reminder", target);
return;
}
pendingReplies.put(target, lead.get());
pendingReplies.compute(target, (t, existing) ->
new ReplyEntry(lead.get(), existing == null ? 0 : existing.nudgeCount()));
startOrCoalesce(lead.get());
}
@@ -232,7 +285,8 @@ public final class ReplyPushLoop {
ticket, target);
return;
}
pendingTickets.put(ticket, new PendingTicket(ticket, lead.get(), failed));
pendingTickets.compute(ticket, (id, existing) ->
new PendingTicket(ticket, lead.get(), failed, existing == null ? 0 : existing.nudgeCount()));
startOrCoalesce(lead.get());
}
@@ -255,28 +309,37 @@ public final class ReplyPushLoop {
return;
}
log.debug("push: starting reminder loop for lead {}", lead);
scheduleNext(lead, 0, 0);
scheduleNext(lead);
}
/** Execute one loop tick — called on the scheduler thread. */
private void tick(String lead, int replyReminderCount, int ticketReminderCount) {
/**
* 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) {
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);
// 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);
scheduleNext(lead);
}
// Re-check after the configured backoff; the lead may become injectable soon.
case WAIT_BUSY -> scheduleNext(lead, replyReminderCount, ticketReminderCount);
case WAIT_BUSY -> scheduleNext(lead);
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore);
}
}
@@ -317,7 +380,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, 0, 0);
scheduleNext(lead);
return;
}
log.debug("push: reminder loop ended for lead {}", lead);
@@ -344,11 +407,28 @@ 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, int nextReplyReminderCount, int nextTicketReminderCount) {
scheduler.schedule(() -> tick(lead, nextReplyReminderCount, nextTicketReminderCount),
private void scheduleNext(String lead) {
scheduler.schedule(() -> tick(lead),
backoffMs, TimeUnit.MILLISECONDS);
}
@@ -13,7 +13,6 @@ 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;
@@ -293,11 +292,6 @@ 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,15 +16,12 @@ 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;
@@ -426,35 +423,6 @@ 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,6 +42,7 @@ 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;
@@ -579,6 +580,83 @@ 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,8 +18,6 @@ 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;
@@ -242,43 +240,6 @@ 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.
+2 -19
View File
@@ -82,25 +82,8 @@
<key>RunAtLoad</key>
<true/>
<!--
CB-600 — read this before assuming ThrottleInterval bounds anything. It paces restarts to at
most one per 10s; it does NOT cap how many times launchd retries. If bridged fails fast on
every start — a bad bridged.yaml, for example auth.mode: token with the token env var unset,
which throws in main() before the daemon ever binds a port — launchd restarts it forever,
once every 10s, until a human intervenes. LaunchAgents have no "give up after N attempts"
primitive, so this is not something a config change here can fix.
That loop stops only two ways: (1) `launchctl unload -w ~/Library/LaunchAgents/dev.ltms.bridged.plist`,
or (2) the underlying cause gets fixed, so the process starts successfully and stays up (no
more exits to restart). scripts/redeploy-bridged.sh does not add a third way — it does not
make bridged self-disable on a config error, on purpose: a fail-fast exit path that
sometimes decides "this is unrecoverable, stop trying" is one more thing that can misfire,
and a wrongly self-disabled daemon needs the exact same manual `launchctl load -w` recovery
this comment already names — so it buys nothing an operator watching for the crash loop
doesn't already have, at the cost of a new way to be silently down. Watch for it with
`launchctl list dev.ltms.bridged` (a high restart count) or by tailing bridged.out for the
same startup error repeating every ~10s.
-->
<!-- Restart on crash, but not in a tight loop if the config is bad (bridged fails fast on a
non-loopback bind without token auth — that is a config error, not a transient one). -->
<key>KeepAlive</key>
<dict>
<key>SuccessfulExit</key>
+1 -66
View File
@@ -79,53 +79,6 @@ running_pid() { pgrep -f "$PATTERN" || true; }
launchd_installed() { [ -f "$LAUNCHD_PLIST" ]; }
launchd_loaded() { launchctl list "$LAUNCHD_LABEL" >/dev/null 2>&1; }
# CB-600: the script computes its own log path from where it sits on disk (REPO, above); the
# plist hard-codes an absolute StandardOutPath. Nothing forced the two to agree — if this script
# were ever run from a checkout other than the one the loaded plist names, launchd would start and
# log the daemon correctly, while every check below (the fresh "bridged listening" line, the
# ERROR-count scan) would read a different, empty or stale file and the script would report a
# clean restart while the daemon crash-loops. Pure and side-effect-free besides `die`/`ok` — reads
# the two paths, resolves them, compares — so it never touches launchd or the daemon and can be
# exercised by sourcing this script (see the SOURCED guard below) without installing the agent.
check_log_path_matches_plist() {
local script_out="$1" plist_path="$2"
local plist_out resolved_out resolved_plist_out
# Checked by exit status, not by emptiness: on a missing file/key PlistBuddy exits nonzero but
# still writes a message ("File Doesn't Exist, Will Create: ...") that command substitution
# would happily capture as if it were the real value — testing only `-z` missed that case.
if ! plist_out="$(/usr/libexec/PlistBuddy -c 'Print :StandardOutPath' "$plist_path" 2>/dev/null)" \
|| [ -z "$plist_out" ]; then
die "launchd agent is loaded but PlistBuddy could not read StandardOutPath from
$plist_path
— cannot verify the daemon logs where this script is about to look. Fix the plist before
redeploying supervised."
fi
resolved_out="$(cd "$(dirname "$script_out")" 2>/dev/null && pwd -P)/$(basename "$script_out")" || true
resolved_plist_out="$(cd "$(dirname "$plist_out")" 2>/dev/null && pwd -P)/$(basename "$plist_out")" || true
if [ -z "$resolved_out" ] || [ -z "$resolved_plist_out" ] || [ "$resolved_out" != "$resolved_plist_out" ]; then
die "log path mismatch — this script reads
$script_out (resolved: ${resolved_out:-<directory does not exist>})
but the loaded plist's StandardOutPath is
$plist_out (resolved: ${resolved_plist_out:-<directory does not exist>})
Under supervision the daemon writes to the PLIST's path, not necessarily this script's — every
post-restart check below (the fresh 'bridged listening' line, the ERROR-count scan) would read
the wrong file and could report a clean restart while the daemon crash-loops. Fix the mismatch
(move this checkout to match the plist, or edit the plist's StandardOutPath/StandardErrorPath)
before redeploying supervised."
fi
ok "log path check: script and plist agree ($resolved_out)"
}
# CB-600: sourceable for testing. When this file is SOURCED (not executed) it stops here — nothing
# below runs — so a test harness can `source` it to call check_log_path_matches_plist (or the
# other pure helpers above) against a throwaway plist fixture without ever reaching the mutating
# flow (build/stop/start) or touching the real daemon or launchd. On a normal `./redeploy-bridged.sh`
# invocation `(return 0 2>/dev/null)` fails (return is illegal at top level of an executed script),
# so this whole block is a no-op and every line below still runs exactly as before.
if (return 0 2>/dev/null); then
return 0
fi
# ---------------------------------------------------------------- report state
say "current state"
@@ -150,9 +103,6 @@ SUPERVISED=0
if launchd_loaded; then
SUPERVISED=1
ok "launchd agent loaded ($LAUNCHD_LABEL) — launchd supervises this daemon"
# CB-600: fail loudly here, before ANY other check runs, if this script and the loaded plist
# would read different log files — every check after this point is worthless otherwise.
check_log_path_matches_plist "$OUT" "$LAUNCHD_PLIST"
else
warn "launchd agent not loaded — this script is the only thing that will restart the daemon."
fi
@@ -272,22 +222,7 @@ say "start"
if [ "$SUPERVISED" = 1 ]; then
echo " supervision is ON: using 'launchctl load' so launchd starts and keeps supervising this"
echo " process, instead of a manual nohup that launchd would know nothing about."
# CB-600: 'launchctl unload -w' above already persisted Disabled=true for this label. A load -w
# that succeeds clears it; a load -w that FAILS leaves the agent both stopped and disabled — worse
# than before this script ran, because a later reboot or login will not bring it back either. One
# retry covers a transient race (e.g. launchd not yet fully done deregistering); if it still fails,
# die with the exact recovery command rather than a bare "failed".
if ! launchctl load -w "$LAUNCHD_PLIST" 2>/dev/null; then
warn "launchctl load failed on the first attempt — retrying once after a short pause"
sleep 2
launchctl load -w "$LAUNCHD_PLIST" || die "launchctl load failed twice.
The agent is now STOPPED and DISABLED — it will NOT come back on its own, not even after a
reboot or login, because 'launchctl unload -w' above persisted Disabled=true and load -w
never got the chance to clear it. Recover with:
launchctl load -w \"$LAUNCHD_PLIST\"
If that still fails, check 'launchctl list $LAUNCHD_LABEL', validate the plist with
'plutil -lint \"$LAUNCHD_PLIST\"', and check $OUT before assuming a retry will succeed."
fi
launchctl load -w "$LAUNCHD_PLIST" || die "launchctl load failed"
else
# Absolute jar path so `ps` names which checkout is running.
( cd "$BRIDGED" && zsh -lc "nohup java -jar '$JAR' >> bridged.out 2>&1 &" )