Compare commits
12 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| fe46311266 | |||
| 27bbd11f06 | |||
| 48d7841fbf | |||
| cc1df11f69 | |||
| fdfd4ac491 | |||
| d5dd5639ae | |||
| 7a120b3256 | |||
| 863d477966 | |||
| a36b7ccd7c | |||
| 32bf324a1e | |||
| 16de9df000 | |||
| 81a0cf4710 |
@@ -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` |
|
||||
|
||||
@@ -275,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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
|
||||
@@ -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<>();
|
||||
|
||||
@@ -35,9 +35,12 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
*/
|
||||
class AmqpReplyInboxRecoveryRaceTest {
|
||||
|
||||
/** Large enough that the (unfixed) unsynchronized sweep's iteration is a real, observable window
|
||||
* a concurrently-started publish can land in — not just a best case, single-entry sprint. */
|
||||
private static final int STALE_PUBLISHES = 100_000;
|
||||
/** Large enough that thousands of entries are still unprocessed by the time the very first one
|
||||
* is observed as failed (see {@code sweepIsHoldingTheLock} below) — that gap is what makes the
|
||||
* head start deterministic instead of a coin flip. 100,000 gave the same guarantee but made the
|
||||
* test far more expensive than the guarantee needs; the ordering no longer depends on a timing
|
||||
* window sized to the full backlog; just to the tail of it. */
|
||||
private static final int STALE_PUBLISHES = 2_000;
|
||||
|
||||
@Test
|
||||
@Timeout(30)
|
||||
@@ -58,6 +61,12 @@ class AmqpReplyInboxRecoveryRaceTest {
|
||||
// pendingByMsgId exactly like publishes whose confirm never arrived before a connection drop.
|
||||
// Virtual threads make this many concurrent blocking publish() calls cheap.
|
||||
CountDownLatch staleStarted = new CountDownLatch(STALE_PUBLISHES);
|
||||
// Counted down by the FIRST stale publish thread to observe its own failure. That can only
|
||||
// happen from inside failPendingPublishesOnRecovery() — nothing else in this test ever
|
||||
// completes a stale Pending exceptionally (no nack/return is simulated for any "stale-*"
|
||||
// msgId) — so seeing it fire is direct, observable proof the sweep is inside its loop, not a
|
||||
// timing guess. It replaces the old fixed Thread.sleep(5) head start.
|
||||
CountDownLatch sweepIsHoldingTheLock = new CountDownLatch(1);
|
||||
for (int i = 0; i < STALE_PUBLISHES; i++) {
|
||||
String msgId = "stale-" + i;
|
||||
Thread.ofVirtual().start(() -> {
|
||||
@@ -65,7 +74,7 @@ class AmqpReplyInboxRecoveryRaceTest {
|
||||
try {
|
||||
inbox.publish("worker-stale", msgId, "x");
|
||||
} catch (IllegalStateException expected) {
|
||||
// resolved (failed by the sweep) — that is exactly what this thread is here for
|
||||
sweepIsHoldingTheLock.countDown();
|
||||
}
|
||||
});
|
||||
}
|
||||
@@ -83,13 +92,28 @@ class AmqpReplyInboxRecoveryRaceTest {
|
||||
}
|
||||
}, "recovery-sweep");
|
||||
sweepThread.start();
|
||||
// A short, deliberate head start: with STALE_PUBLISHES this large, the (unfixed) sweep's own
|
||||
// iteration takes several milliseconds, so this guarantees the sweep has already begun —
|
||||
// and, once guarded, is already holding publishChannelLock — before "fresh" attempts to
|
||||
// register. Without this head start, "fresh" sometimes wins the race for the lock and
|
||||
// registers before the sweep even starts, which is the accepted "already in flight when
|
||||
// recovery fires" case (correctly failed either way) rather than the bug under test.
|
||||
Thread.sleep(5);
|
||||
|
||||
// Deterministic head start: block until the sweep has actually failed one of the stale
|
||||
// publishes. failPendingPublishesOnRecovery() (once guarded, as it is on main) holds
|
||||
// publishChannelLock for its ENTIRE loop, not just per entry — so this failure proves the
|
||||
// sweep is, at this instant, still holding that lock. With STALE_PUBLISHES this large, the
|
||||
// remaining ~1,999 entries give an enormous margin between "first failure observed" and "sweep
|
||||
// releases the lock": there is no window left for "fresh" to slip in before the sweep starts,
|
||||
// or to win the lock ahead of it — see the case-2 note below. This also means Case 1 (the sweep
|
||||
// is already inside its loop, holding the lock, when "fresh" tries to register) is now
|
||||
// guaranteed by construction rather than merely likely under a fixed sleep.
|
||||
assertTrue(sweepIsHoldingTheLock.await(20, TimeUnit.SECONDS),
|
||||
"the sweep never failed a single stale publish — it may not have started");
|
||||
|
||||
// Case 2 ("fresh" wins publishChannelLock before the sweep even starts, so it genuinely
|
||||
// published on the stale channel and the sweep correctly fails it) is impossible by
|
||||
// construction in this test: freshThread.start() below is reached only after
|
||||
// sweepIsHoldingTheLock has counted down, which can only happen once
|
||||
// failPendingPublishesOnRecovery() is already running and has already failed a stale entry.
|
||||
// There is no code path that lets "fresh" start before the sweep starts. That case is real
|
||||
// and correct production behaviour (see AmqpReplyInbox#failPendingPublishesOnRecovery's
|
||||
// javadoc), it is just not reachable from this deterministic ordering, so it does not need a
|
||||
// separate assertion here.
|
||||
|
||||
// This is the exact interleaving CB-528's follow-up describes: "the still-running recovery
|
||||
// sweep" racing a publish that registers while it is mid-flight.
|
||||
@@ -109,8 +133,8 @@ class AmqpReplyInboxRecoveryRaceTest {
|
||||
// Simulate the broker's real confirm for "fresh" now that the sweep is done, so a correct
|
||||
// implementation's publish() returns normally instead of idling out CONFIRM_TIMEOUT_MS. Poll
|
||||
// for the registration rather than checking once: freshThread may still be contending for
|
||||
// publishChannelLock (behind the 20,000 stale threads' own lock acquisitions) even though the
|
||||
// sweep itself has already finished.
|
||||
// publishChannelLock (behind the sweep's own hold on it, and possibly other stale threads
|
||||
// still unwinding) even though the sweep itself has already finished.
|
||||
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(9);
|
||||
int idx = -1;
|
||||
while (idx < 0 && System.nanoTime() < deadline) {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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();
|
||||
|
||||
Reference in New Issue
Block a user