Compare commits
8 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| bb750cdba3 | |||
| d56c77b368 | |||
| 83e2ff06cf | |||
| f5deaafd06 | |||
| fe46311266 | |||
| 27bbd11f06 | |||
| 48d7841fbf | |||
| cc1df11f69 |
@@ -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
|
||||
|
||||
@@ -7,6 +7,7 @@ import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
|
||||
import dev.ltms.bridged.msg.AmqpReplyInbox;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
import dev.ltms.bridged.placement.PlacementPolicies;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -986,7 +987,14 @@ public record BridgedConfig(
|
||||
warnUnknownTopLevelKeys(yaml, path);
|
||||
rejectDuplicateMemberSlots(yaml);
|
||||
rejectNegativeMaxLoad(yaml);
|
||||
rejectUnknownKind(yaml);
|
||||
rejectUnknownAuthMode(yaml);
|
||||
rejectUnknownPlacement(yaml);
|
||||
BridgedConfig cfg = YAML.readValue(yaml, BridgedConfig.class);
|
||||
// CB-606: validated here, eagerly, using PlacementPolicies.fromName as the single source
|
||||
// of truth — not lazily at first spawn (see CompositePeerLauncher's placementPolicy
|
||||
// Supplier), where a bad name would still start a daemon that looks healthy.
|
||||
rejectUnknownPlacementPolicy(cfg.placement());
|
||||
return cfg.withDefaults();
|
||||
} catch (IOException e) {
|
||||
throw new UncheckedIOException("cannot read bridged config at " + path, e);
|
||||
@@ -1280,6 +1288,162 @@ 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.");
|
||||
}
|
||||
}
|
||||
|
||||
/** The auth modes this build understands — {@link Auth#mode()}'s only valid values. */
|
||||
private static final Set<String> KNOWN_AUTH_MODES = Set.of(Auth.MODE_LOOPBACK_TRUST, Auth.MODE_TOKEN);
|
||||
|
||||
/**
|
||||
* Reject an {@code auth.mode} that is not one of {@link #KNOWN_AUTH_MODES} (CB-606), naming the
|
||||
* value and the accepted set.
|
||||
*
|
||||
* <p>{@link Auth}'s compact constructor only lower-cases {@code mode}, and
|
||||
* {@link Auth#tokenMode()} only compares the result against {@code MODE_TOKEN} — anything else,
|
||||
* including a typo like {@code toekn}, silently behaves as {@code loopback-trust}. That fallback
|
||||
* is otherwise checked only by {@link #validateAuthExposure()}, and only when the bind is
|
||||
* non-loopback: on a loopback bind (the common case) the typo is invisible end to end — the
|
||||
* daemon starts cleanly and authenticates nobody while the operator believes {@code token} mode
|
||||
* is active. Refuse it here, unconditionally, at config load, rather than let it hide behind the
|
||||
* bind check.
|
||||
*
|
||||
* @param yaml the raw config text
|
||||
* @throws IllegalStateException when {@code auth.mode} is a non-blank value not in
|
||||
* {@link #KNOWN_AUTH_MODES} (case-insensitive)
|
||||
*/
|
||||
static void rejectUnknownAuthMode(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("auth") instanceof Map<?, ?> auth)) {
|
||||
return;
|
||||
}
|
||||
if (!(auth.get("mode") instanceof String mode) || mode.isBlank()
|
||||
|| KNOWN_AUTH_MODES.contains(mode.toLowerCase())) {
|
||||
return;
|
||||
}
|
||||
throw new IllegalStateException("refusing to start: auth.mode=" + mode
|
||||
+ " is not recognized — accepted values are "
|
||||
+ String.join(", ", KNOWN_AUTH_MODES.stream().sorted().toList())
|
||||
+ " (case-insensitive); an unrecognized mode would otherwise silently fall back to"
|
||||
+ " loopback-trust, which authenticates nobody.");
|
||||
}
|
||||
|
||||
/** The per-profile placements this build understands — {@link Profile#placement()}'s only valid values. */
|
||||
private static final Set<String> KNOWN_PLACEMENTS = Set.of("tab", "pane");
|
||||
|
||||
/**
|
||||
* Reject a profile whose {@code placement:} is not one of {@link #KNOWN_PLACEMENTS} (CB-606),
|
||||
* naming the profile, the value it set, and the accepted set.
|
||||
*
|
||||
* <p>{@link Profile}'s compact constructor only lower-cases {@code placement}, and
|
||||
* {@link Profile#tabPlacement()} only compares the result against {@code "tab"} — anything else,
|
||||
* including a typo like {@code tabb}, silently falls back to the legacy pane placement with no
|
||||
* signal anywhere.
|
||||
*
|
||||
* @param yaml the raw config text
|
||||
* @throws IllegalStateException when any profile's {@code placement} is a non-blank value not in
|
||||
* {@link #KNOWN_PLACEMENTS} (case-insensitive)
|
||||
*/
|
||||
static void rejectUnknownPlacement(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("placement") instanceof String pl && !pl.isBlank()
|
||||
&& !KNOWN_PLACEMENTS.contains(pl.toLowerCase()))
|
||||
.map(e -> String.valueOf(e.getKey()) + "=" + ((Map<?, ?>) e.getValue()).get("placement"))
|
||||
.sorted()
|
||||
.toList();
|
||||
if (!bad.isEmpty()) {
|
||||
throw new IllegalStateException("refusing to start: profile(s) [" + String.join(", ", bad)
|
||||
+ "] set an unrecognized placement — accepted values are "
|
||||
+ String.join(", ", KNOWN_PLACEMENTS.stream().sorted().toList())
|
||||
+ " (case-insensitive); an unrecognized placement would otherwise fall back to"
|
||||
+ " legacy pane placement with no signal anywhere.");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Reject a top-level {@code placement:} policy name {@link PlacementPolicies#fromName} does not
|
||||
* recognize (CB-606), at config load rather than lazily at first spawn.
|
||||
*
|
||||
* <p>{@code CompositePeerLauncher} only calls {@link PlacementPolicies#fromName} per spawn,
|
||||
* through a {@code Supplier} that re-reads live config (CB-559, so a hot-reloaded placement
|
||||
* policy takes effect without a restart) — so a bad name still starts a daemon that looks
|
||||
* healthy and fails only the first time something spawns without naming a profile. Every other
|
||||
* field this class validates fails here, at load; this one gets the same treatment, calling
|
||||
* {@link PlacementPolicies#fromName} itself as the single source of truth for what is valid
|
||||
* rather than duplicating its accepted set.
|
||||
*
|
||||
* @param placement the raw, possibly null/blank {@code placement} value as parsed (before
|
||||
* {@link #withDefaults()} runs); {@code fromName} itself treats null/blank as
|
||||
* {@code fixed}, so this call changes no default
|
||||
* @throws IllegalStateException when {@code placement} is a name {@link PlacementPolicies} does
|
||||
* not recognize
|
||||
*/
|
||||
private static void rejectUnknownPlacementPolicy(String placement) {
|
||||
try {
|
||||
PlacementPolicies.fromName(placement);
|
||||
} catch (IllegalArgumentException e) {
|
||||
throw new IllegalStateException("refusing to start: " + e.getMessage(), e);
|
||||
}
|
||||
}
|
||||
|
||||
static List<String> unknownTopLevelKeys(String yaml) {
|
||||
Map<?, ?> raw;
|
||||
try {
|
||||
|
||||
@@ -536,6 +536,14 @@ public final class MessageService {
|
||||
// Only the fresh owner tears down the shared turn; a duplicate must leave it open.
|
||||
if (ticket.fresh()) {
|
||||
rendezvous.closeAsk(ticket.turnId());
|
||||
// CB-582: tear the push loop's copy down at the same point, not only on the three
|
||||
// paths that call clearAsyncQuestion. The answer future can complete exceptionally
|
||||
// (ExecutionException) or the thread be interrupted, and both leave this method by
|
||||
// throwing — the question would stay pending forever, keep being named in nudges
|
||||
// until its own cap, and never be removed from the map. Already-closed is a no-op.
|
||||
if (pushLoop != null) {
|
||||
pushLoop.questionClosed(ticket.turnId());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -1040,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");
|
||||
@@ -1369,6 +1407,89 @@ class BridgedConfigTest {
|
||||
assertTrue(e.getMessage().contains("maxLoad"), "error names the key: " + e.getMessage());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-606: an unrecognized {@code auth.mode} used to silently fall back to
|
||||
* {@code loopback-trust} — {@link BridgedConfig.Auth#tokenMode()} only checked equality
|
||||
* against {@code "token"}. On a loopback bind {@link BridgedConfig#validateAuthExposure()}
|
||||
* never runs (it only fires for a non-loopback bind), so the typo was completely invisible:
|
||||
* the daemon started cleanly and authenticated nobody while the operator believed token mode
|
||||
* was active.
|
||||
*/
|
||||
@Test
|
||||
void unknownAuthModeIsRefusedAtLoadNamingTheValueAndTheAcceptedSet(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("auth-mode-typo.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
auth:
|
||||
mode: toekn
|
||||
""");
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("toekn"), "error names the bad value: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("loopback-trust") && e.getMessage().contains("token"),
|
||||
"error names the accepted set: " + e.getMessage());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-606: an unrecognized per-profile {@code placement:} used to silently fall back to legacy
|
||||
* pane placement — {@link BridgedConfig.Profile#tabPlacement()} only checked equality against
|
||||
* {@code "tab"}.
|
||||
*/
|
||||
@Test
|
||||
void unknownProfilePlacementIsRefusedAtLoadNamingTheProfileAndTheAcceptedSet(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("placement-typo.yaml");
|
||||
Files.writeString(f, """
|
||||
profiles:
|
||||
gx10:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
placement: tabb
|
||||
""");
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("gx10"), "error names the profile: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("tabb"), "error names the bad value: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("tab") && e.getMessage().contains("pane"),
|
||||
"error names the accepted set: " + e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void absentProfilePlacementDefaultsToTab(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("placement-absent.yaml");
|
||||
Files.writeString(f, """
|
||||
profiles:
|
||||
gx10:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
""");
|
||||
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
BridgedConfig.Profile w = cfg.profiles().get("gx10");
|
||||
assertTrue(w.tabPlacement(), "an absent placement must keep defaulting to tab");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-606: the top-level {@code placement:} policy name WAS validated, but only lazily, by
|
||||
* {@code PlacementPolicies.fromName} through {@code CompositePeerLauncher}'s per-spawn
|
||||
* {@code Supplier} — so a bad name still started a daemon that looked healthy and failed only
|
||||
* the first time something spawned without naming a profile. This must now fail at load.
|
||||
*/
|
||||
@Test
|
||||
void unknownTopLevelPlacementPolicyIsRefusedAtLoadNotLazilyAtFirstSpawn(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("placement-policy-typo.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
placement: weightd
|
||||
""");
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("weightd"), "error names the bad value: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("fixed") && e.getMessage().contains("round-robin")
|
||||
&& e.getMessage().contains("weighted"),
|
||||
"error names the accepted set: " + e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void subscriptionFlagBindsAndDefaultsFalse(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("subscription.yaml");
|
||||
|
||||
@@ -939,6 +939,45 @@ class MessageServiceTest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void anAskThatLeavesByThrowingStillClosesItsQuestion() throws Exception {
|
||||
// CB-582 follow-up. ask() calls clearAsyncQuestion on three paths — no-waiter, timed out,
|
||||
// and (from answer()) answered — but it can also leave by *throwing*: an interrupt while
|
||||
// blocked on the answer, or an ExecutionException from the answer future. Those paths run
|
||||
// only the finally block, so before the fix the push loop kept the question pending for
|
||||
// good: named in every nudge until its own cap, then never removed from the map at all.
|
||||
PrimaryRegistry registry = new PrimaryRegistry(null);
|
||||
registry.recordDelegation(T, LEAD);
|
||||
var scheduler = java.util.concurrent.Executors.newSingleThreadScheduledExecutor();
|
||||
// A backoff far longer than the test: the schedule is started but no tick ever fires, so
|
||||
// decide() is read directly and nothing here depends on timing.
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, new AgentControl(new FakeHerdr()), inbox,
|
||||
scheduler, 5, 60_000);
|
||||
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop);
|
||||
try {
|
||||
String ticket = service.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
Thread asker = new Thread(() -> assertThrows(IllegalStateException.class,
|
||||
() -> service.ask(T, "which config file?", 30_000)));
|
||||
asker.start();
|
||||
awaitTicketPhaseOn(service, ticket, MessageService.Phase.ASKING);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, pushLoop.decide(LEAD, 0, 0, 0),
|
||||
"the open question should be the one thing keeping this lead's schedule alive");
|
||||
|
||||
asker.interrupt();
|
||||
asker.join(5000);
|
||||
assertFalse(asker.isAlive(), "the interrupted ask should have left ask() by throwing");
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.STOP, pushLoop.decide(LEAD, 0, 0, 0),
|
||||
"an ask that threw must still close its question, or the loop nudges about it for good");
|
||||
} finally {
|
||||
service.close();
|
||||
scheduler.shutdownNow();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFleetWithNoPushLoopConfiguredBehavesExactlyAsToday() throws Exception {
|
||||
// `messages` (the shared field) uses the no-pushLoop constructor — poll() must not throw,
|
||||
|
||||
@@ -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