Compare commits

...

16 Commits

Author SHA1 Message Date
Dai Ha fe46311266 CB-604: reject an unrecognized profile kind at config load
CI / contract (pull_request) Successful in 53s
CI / build (pull_request) Successful in 1m10s
2026-08-16 18:37:21 +02:00
Dai Ha 27bbd11f06 Merge CB-584: carry the agent session id on a failed ticket's detail
CI / contract (push) Successful in 47s
CI / build (push) Successful in 1m23s
Most of issue #65 was already shipped in 5d5b3bd - MemberSession
records agentSessionId, PeerHandle.agentSessionId() has no default, the
roster exposes it, and bridge_spawn accepts sessionName and
resumeSessionId. That commit left one item for follow-up: carrying the
id on ReleaseDetail.

This does that item. CB-578 stage C already lets a lead re-dispatch onto
the same worktree after a failure; without the session id that is a cold
start. With it, the work and the thread both survive.

Also updates the bridge_spawn and bridge_list rows in the CLAUDE.md
intent table, which the earlier commit missed.

Verified here: trial merge onto main builds 830 tests BUILD SUCCESS,
unpiped. Confirmed against the code that items 1-4 really were already
on main, so the scope-down is correct rather than work skipped.

Closes #65
2026-08-16 18:29:24 +02:00
Dai Ha 48d7841fbf CB-589: document that weighted placement is not cheapest-first
CI / contract (push) Successful in 50s
CI / build (push) Successful in 1m43s
The weight ratio does not express a preference order. weighted spreads
spawns across every profile with a free slot, so paid spawns happen
while the free box is idle - and a profile at maxLoad freezes its score,
so it can lose the next pick after a slot frees.

The real fix is a cost-first policy (CB-589). This documents the
workaround and its trap next to the key, because bridged.yaml is
gitignored: a fresh host starts without the workaround and quietly pays,
with nothing to tell the operator why.

Comments only. BridgedConfigTest: 85 tests, BUILD SUCCESS.
2026-08-16 18:27:19 +02:00
Dai Ha cc1df11f69 CB-584: carry agentSessionId on a failed ticket's ReleaseDetail
CI / contract (pull_request) Successful in 1m6s
CI / build (pull_request) Successful in 1m42s
Closes the one piece issue #65 deliberately left out of 5d5b3bd: a
released member's agentSessionId now rides alongside worktree, branch
and snapshotRef on ReleaseDetail, and Bridged's onRelease handler
names it in the abandon reason, so a lead can resume the member's
conversation instead of only re-dispatching a fresh one onto the same
files.

Also updates the CLAUDE.md bridge_spawn/bridge_list table row, which
5d5b3bd shipped the sessionName/resumeSessionId/agentSessionId surface
for but never updated.
2026-08-16 18:25:24 +02:00
Dai Ha fdfd4ac491 Merge CB-600: make installing the launchd agent safe
CI / contract (push) Successful in 51s
CI / build (push) Successful in 1m42s
Three gaps that only bite once the agent is loaded, plus one wrong
comment.

The script computed its log path from its own location while the plist
hard-codes one. Run from a different checkout, every post-restart check
would read the wrong file and report a clean restart while the daemon
crash-looped. It now compares the two and fails, not warns.

A failed 'launchctl load' after a successful 'unload -w' left the agent
stopped AND persistently disabled - worse than before the redeploy. It
now retries once, then dies naming the exact recovery command.

The plist now says plainly that ThrottleInterval paces restarts but does
not bound them, and what actually stops the loop.

Verified here: ran the script with --check from the merged tree and it
behaves exactly as before, so the unsupervised path - my only restart
route - is intact. Exercised the log-path check against match, mismatch
and missing-plist fixtures using a truncated copy with no mutating code
in it: ok/1/1. 829 tests BUILD SUCCESS.

Closes #91
2026-08-16 18:17:11 +02:00
Dai Ha d5dd5639ae Merge CB-602: guard against a config key that never reaches the example
bridged.yaml is gitignored, so bridged.example.yaml is the only
committed description of the config schema. Two tests already covered
example -> code; nothing covered code -> example, so a brand-new key
could ship undocumented and no test would notice.

A new test compares BridgedConfig.KNOWN_TOP_LEVEL_KEYS against the
example scanned as TEXT, so a key documented only as a comment counts as
documented. That is what makes the guard correct rather than annoying:
most of the example is commented on purpose.

Verified here: added an undocumented key and watched the test fail with
an actionable message naming it; then documented that key as a comment
only and watched it pass. Probe reverted, tree clean.

Closes #96
2026-08-16 18:17:11 +02:00
Dai Ha 7a120b3256 Merge CB-598: per-item reminder counts so backoff-window work is never orphaned
CI / contract (push) Successful in 45s
CI / build (push) Successful in 1m39s
The reminder count was one counter per lead per source, carried forward
across ticks. A counter carried forward has no memory of which item it
counted, so work arriving during the backoff window inherited an
already-capped count and was never named in a nudge.

tick() now recomputes each source's count fresh from the minimum count
among the items actually pending, tracked per item. A fresh item keeps
its source eligible; an older capped item still rides along in the text
without spending more budget. decide() is unchanged.

Verified here: read the diff; the bumped set is exactly the set named in
the nudge, and the empty early-return skips the bump. Trial merge onto
main builds 826 tests BUILD SUCCESS, unpiped.

Closes #87
2026-08-16 18:13:29 +02:00
Dai Ha 863d477966 CB-603: make FakeHerdr.calls thread-safe
Background loops call the fake from their own scheduler threads while a
test polls called() from the test thread. The list was a plain
ArrayList, so a nudge landing mid-stream threw
ConcurrentModificationException out of called().

It surfaced while I was verifying CB-598, which nudges more often, but
the race is on main today and is unrelated to that change.

824 tests, BUILD SUCCESS.
2026-08-16 18:12:23 +02:00
Dai Ha cec48832be CB-600: make it safe to install the launchd agent
CI / build (pull_request) Failing after 1m21s
CI / contract (pull_request) Successful in 1m26s
- redeploy-bridged.sh now refuses (not warns) a supervised restart when
  its computed log path disagrees with the loaded plist's StandardOutPath
  — otherwise every post-restart check reads the wrong file and can
  report a clean restart while the daemon crash-loops. The check is a
  pure, testable function; the script gained a source-for-test guard so
  it can be exercised without installing the agent or touching launchd.
- a failed 'launchctl load' after a successful 'unload' now retries once
  and, on ultimate failure, tells the operator the agent is stopped AND
  disabled plus the exact recovery command, instead of leaving that
  silently worse than the pre-redeploy state.
- the plist documents honestly that the crash loop launchd retries is
  unbounded (ThrottleInterval only paces it), and what actually stops it.
- fixed the requiredSecretEnvVars javadoc: the auth.tokenEnv startup
  throw is ~370 lines below its call site, not a few lines above it, and
  only fires in auth.mode: token.
2026-08-16 18:08:58 +02:00
Dai Ha a36b7ccd7c Merge CB-601: make the recovery-race test deterministic
CI / build (push) Successful in 1m5s
CI / contract (push) Successful in 1m7s
The test asserted one of two interleavings that are both correct, and
steered toward it with a 5 ms Thread.sleep. Under load the other
interleaving happened and main went red on a correct implementation.

The head start is now a latch counted down from inside the sweep's
guarded loop, so the ordering is guaranteed, not likely. Test file only;
the production guard is unchanged.

Verified here: 822 tests BUILD SUCCESS; 10/10 passes while a full clean
install ran in parallel; 5/5 failures with the guard removed, so the test
still catches the bug it exists for.

Closes #95
2026-08-16 18:08:18 +02:00
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 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
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
16 changed files with 617 additions and 42 deletions
+2 -2
View File
@@ -105,8 +105,8 @@ the merge — and merging on a reviewer's word is delegating it by proxy.
|---|---|
| Confirm your own role | `bridge_whoami` |
| See backends available | `bridge_profiles` |
| Start a member | `bridge_spawn{role?, profile?, cwd?, worktree?, ticket?}` → `sessionId` + `paneId` |
| See the fleet | `bridge_list` → `leads` (your peers) + `members` · one peer's state: `bridge_status{sessionId}` |
| Start a member | `bridge_spawn{role?, profile?, cwd?, worktree?, ticket?, sessionName?, resumeSessionId?}` → `sessionId` + `paneId` |
| See the fleet | `bridge_list` → `leads` (your peers) + `members` (each carries `agentSessionId` when its backend knows one) · one peer's state: `bridge_status{sessionId}` |
| Delegate (blocking) | `bridge_send{sessionId, content}` |
| Delegate (long task) | `bridge_send{sessionId, content, wait:false}` → ticket → `bridge_poll{ticket}` |
| Answer a member's `bridge_ask` | `bridge_send{turnId, content}` — **not** `sessionId` |
+55 -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}).
@@ -262,6 +275,26 @@ profiles:
# argv: ["opencode"]
# How an unqualified spawn chooses a profile: fixed (default, reproduces pre-CB-518 behaviour),
# round-robin, or weighted. Omitting this key is a strict no-op for existing configs.
#
# `weighted` IS NOT "cheapest first" — read this before you set weights (CB-589).
# It is smooth weighted round-robin: it spreads spawns across EVERY profile that has a free slot,
# in weight ratio. It has no idea which profile costs money. So with local:10 / paid:2 you do not
# get "use local, overflow to paid" — you get roughly one spawn in six going to the paid profile
# while the local box still has a free slot.
#
# There is a sharper second effect. The policy's running score map lives for the daemon's whole
# life. While a profile is at maxLoad it is filtered out and its score FREEZES, so the paid
# profiles keep accumulating against it. When the local slot frees up it returns with a stale
# score and can LOSE the next pick — a paid spawn while the free box sits idle.
#
# Until a real cost-first policy exists, the workaround is to make the ratio decisive rather than
# proportional: give the free profile a weight so large that it wins every pick it is eligible
# for, and paid profiles only ever take genuine overflow. On this host that is local weight 100
# against paid weights of ~1.
#
# The gotcha with that workaround: it expresses a PREFERENCE ORDER through a RATIO knob. Add a
# future profile at weight 150 and it silently outranks the free box, with nothing to warn you.
# Re-check the weights whenever you add a profile.
placement: weighted
# How long a credential sits out after a BACKEND_EXHAUSTED classification (CB-578 stage B), in
@@ -286,6 +319,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 +406,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 +467,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).
@@ -441,6 +441,11 @@ public final class Bridged {
reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch()
+ " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none");
}
// CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the
// member's conversation instead of only re-dispatching a fresh one onto the same files.
if (detail.agentSessionId() != null) {
reason += " agentSessionId=" + detail.agentSessionId();
}
messages.abandon(detail.terminalId(), reason);
replyInbox.release(detail.terminalId());
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
@@ -555,7 +560,10 @@ 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, a few lines above this method's call site.
* 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.
*
* <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.
@@ -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");
@@ -982,6 +986,7 @@ public record BridgedConfig(
warnUnknownTopLevelKeys(yaml, path);
rejectDuplicateMemberSlots(yaml);
rejectNegativeMaxLoad(yaml);
rejectUnknownKind(yaml);
BridgedConfig cfg = YAML.readValue(yaml, BridgedConfig.class);
return cfg.withDefaults();
} catch (IOException e) {
@@ -1276,6 +1281,53 @@ public record BridgedConfig(
}
}
/** The peer kinds this build has an adapter for — {@link Profile#kind()}'s only valid values. */
private static final Set<String> KNOWN_KINDS = Set.of(Profile.KIND_CLAUDE_CODE, Profile.KIND_OPENCODE);
/**
* Reject a profile whose {@code kind:} is not one of {@link #KNOWN_KINDS} (CB-604), naming the
* profile, the value it set, and the accepted set.
*
* <p>{@link Profile}'s compact constructor only lower-cases {@code kind} and compares it against
* {@code KIND_OPENCODE} — anything else, including a typo like {@code opencod}, silently falls
* into the claude-code bucket ({@link dev.ltms.bridged.member.CompositePeerLauncher} routes by
* exact adapter claim, not by membership in a known set). With {@code argv:} also unset, the argv
* default special-cases only the exact string {@code "claude-code"}, so the launch command falls
* back to {@code List.of(kind)} — the daemon then tries to run a program literally named after the
* typo. {@code CompositePeerLauncher}'s constructor already treats a profile claimed by two
* adapters as fatal (CB-402); an unrecognized kind is the same class of adapter-routing mistake
* and gets the same treatment here, at config load, rather than surfacing later as a failed spawn.
*
* @param yaml the raw config text
* @throws IllegalStateException when any profile's {@code kind} is a non-blank value not in
* {@link #KNOWN_KINDS} (case-insensitive)
*/
static void rejectUnknownKind(String yaml) {
Map<?, ?> raw;
try {
raw = YAML.readValue(yaml, Map.class);
} catch (IOException | IllegalArgumentException e) {
return; // a malformed file is reported by the real parse, not here
}
if (raw == null || !(raw.get("profiles") instanceof Map<?, ?> profiles)) {
return;
}
List<String> bad = profiles.entrySet().stream()
.filter(e -> e.getValue() instanceof Map<?, ?> p
&& p.get("kind") instanceof String k && !k.isBlank()
&& !KNOWN_KINDS.contains(k.toLowerCase()))
.map(e -> String.valueOf(e.getKey()) + "=" + ((Map<?, ?>) e.getValue()).get("kind"))
.sorted()
.toList();
if (!bad.isEmpty()) {
throw new IllegalStateException("refusing to start: profile(s) [" + String.join(", ", bad)
+ "] set an unrecognized kind — accepted values are "
+ String.join(", ", KNOWN_KINDS.stream().sorted().toList())
+ " (case-insensitive); an unrecognized kind would otherwise fall back to the"
+ " claude-code adapter and try to launch a program named after the typo.");
}
}
static List<String> unknownTopLevelKeys(String yaml) {
Map<?, ?> raw;
try {
@@ -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,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,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) {
@@ -293,8 +293,10 @@ public final class SessionManager implements TurnListener {
// blocked caller fails fast with a real reason instead of sitting on a rendezvous
// nothing will ever resolve. CB-578 stage C: carry the worktree/branch/snapshot ref
// too, so a failed ticket's detail can point a lead at the same tree to re-dispatch.
// CB-584 (issue #65 criterion 5): carry agentSessionId alongside them, so a lead can
// also resume the member's conversation, not just re-dispatch onto its files.
notifyReleased(new ReleaseDetail(removed.terminalId(), removed.worktree(),
removed.branch(), snapshotRef));
removed.branch(), snapshotRef, removed.agentSessionId()));
}
}
// CB-581: the pane must always stop, even if the dirty check above threw. A session removed
@@ -348,7 +350,8 @@ public final class SessionManager implements TurnListener {
* are {@code null} for a shared-tree session; {@code snapshotRef} is {@code null} unless this
* release snapshotted a dirty worktree into {@code refs/wip/<branch>}.
*/
public record ReleaseDetail(String terminalId, String worktreePath, String branch, String snapshotRef) {
public record ReleaseDetail(String terminalId, String worktreePath, String branch, String snapshotRef,
String agentSessionId) {
}
/**
@@ -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.*;
@@ -1039,6 +1040,44 @@ class BridgedConfigTest {
"an opencode worker with no argv defaults to the opencode binary, never claude");
}
/**
* CB-604: an unrecognized {@code kind:} used to silently fall into the claude-code bucket — not
* matching {@code "opencode"} was the only check. With {@code argv:} also unset that meant the
* daemon tried to launch a program literally named after the typo.
*/
@Test
void unknownKindIsRefusedAtLoadNamingTheValueAndTheAcceptedSet(@TempDir Path dir) throws Exception {
Path f = dir.resolve("kind-typo.yaml");
Files.writeString(f, """
profiles:
gemini:
kind: opencod
model: google/gemini-2.5-pro
""");
IllegalStateException e = assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
assertTrue(e.getMessage().contains("gemini"), "error names the profile: " + e.getMessage());
assertTrue(e.getMessage().contains("opencod"), "error names the bad value: " + e.getMessage());
assertTrue(e.getMessage().contains("claude-code") && e.getMessage().contains("opencode"),
"error names the accepted set: " + e.getMessage());
}
@Test
void blankKindStillDefaultsToClaudeCode(@TempDir Path dir) throws Exception {
Path f = dir.resolve("kind-blank.yaml");
Files.writeString(f, """
profiles:
gx10:
baseUrl: http://gx10.gw:8000
kind: ""
argv: ["ccs", "gx10"]
""");
BridgedConfig cfg = BridgedConfig.load(f);
assertEquals(BridgedConfig.Profile.KIND_CLAUDE_CODE, cfg.profiles().get("gx10").kind(),
"a blank kind: is documented to behave exactly like an absent one");
}
@Test
void authDefaultsToLoopbackTrustSoExistingConfigsBehaveAsBefore(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-auth-block.yaml");
@@ -1189,6 +1228,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");
@@ -7,6 +7,7 @@ import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CopyOnWriteArrayList;
/**
* Recording fake {@link HerdrClient} for unit/acceptance tests. Returns canned frames
@@ -22,7 +23,13 @@ public final class FakeHerdr implements HerdrClient {
public static final long WORKER_PID = 4242;
private final ObjectMapper mapper = new ObjectMapper();
public final List<Call> calls = new ArrayList<>();
/**
* Thread-safe on purpose. Background loops — {@link dev.ltms.bridged.msg.ReplyPushLoop} and the
* lead heartbeat — call this fake from their own scheduler threads while a test polls
* {@link #called} from the test thread. A plain {@code ArrayList} threw
* {@code ConcurrentModificationException} out of {@code called()} when a nudge landed mid-stream.
*/
public final List<Call> calls = new CopyOnWriteArrayList<>();
private boolean healthy = true;
private final List<String> extraWorkspaces = new ArrayList<>();
private final List<String> extraAgents = new ArrayList<>();
@@ -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,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,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.
@@ -422,6 +422,28 @@ class WorktreeSessionManagerTest {
"acceptance criterion 6: a failed ticket's detail must carry the snapshot ref");
}
@Test
void releaseNotifiesTheListenerWithTheAgentSessionId() {
// CB-584 (issue #65 criterion 5): a failed ticket's detail must also name the agent
// session, alongside worktree/branch/snapshot, so a lead can resume the conversation
// rather than only re-dispatch a fresh member onto the same files.
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
.withDirty(true);
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
java.util.List<SessionManager.ReleaseDetail> released = new java.util.concurrent.CopyOnWriteArrayList<>();
sessions.onRelease(released::add);
MemberSession s = sessions.acquire("ltms-local", MemberRole.DEV, null, "/caller/proj", null,
new WorktreeRequest("cb-584-e", null), "cb-584-session", null);
assertNotNull(s.agentSessionId(), "a named session must mint an agent session id to assert on");
sessions.release(s.paneId());
assertEquals(1, released.size());
SessionManager.ReleaseDetail detail = released.getFirst();
assertEquals(s.agentSessionId(), detail.agentSessionId());
}
@Test
void aFailingSnapshotStillPreservesTheWorktreeStopsThePaneAndNotifies() {
FakeHerdr herdr = new FakeHerdr();
+19 -2
View File
@@ -82,8 +82,25 @@
<key>RunAtLoad</key>
<true/>
<!-- 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). -->
<!--
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.
-->
<key>KeepAlive</key>
<dict>
<key>SuccessfulExit</key>
+66 -1
View File
@@ -79,6 +79,53 @@ 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"
@@ -103,6 +150,9 @@ 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
@@ -222,7 +272,22 @@ 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."
launchctl load -w "$LAUNCHD_PLIST" || die "launchctl load failed"
# 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
else
# Absolute jar path so `ps` names which checkout is running.
( cd "$BRIDGED" && zsh -lc "nohup java -jar '$JAR' >> bridged.out 2>&1 &" )