Merge remote-tracking branch 'origin/main' into fix-charter
# Conflicts: # bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java
This commit is contained in:
@@ -194,7 +194,7 @@ class BridgedConfigTest {
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
tab: "lead: opus"
|
||||
""");
|
||||
|
||||
BridgedConfig.Leader lead = BridgedConfig.load(f).fleet().leaders().get("opus");
|
||||
@@ -212,7 +212,7 @@ class BridgedConfigTest {
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
tab: "drive: opus"
|
||||
tabPrefix: "drive:"
|
||||
scanIntervalSeconds: 30
|
||||
""");
|
||||
@@ -224,7 +224,9 @@ class BridgedConfigTest {
|
||||
|
||||
/**
|
||||
* The pane no longer has to exist before the daemon does (CB-557): a lead naming a profile may
|
||||
* be launched, while one that names only a terminal is recognised and never created.
|
||||
* be launched, while one that names no profile is recognised and never created. Either way it
|
||||
* still needs its own {@code tab:} (CB-579) — that part is unconditional, see
|
||||
* {@link #aLeadWithNoTabRefusesToStart}.
|
||||
*/
|
||||
@Test
|
||||
void aLeadIsCreatableOnlyWhenItNamesAProfile(@TempDir Path dir) throws Exception {
|
||||
@@ -239,8 +241,9 @@ class BridgedConfigTest {
|
||||
leaders:
|
||||
launched:
|
||||
profile: opus
|
||||
tab: "lead: launched"
|
||||
pinned:
|
||||
terminal: term_opus
|
||||
tab: "lead: pinned"
|
||||
""");
|
||||
|
||||
var leaders = BridgedConfig.load(f).fleet().leaders();
|
||||
@@ -249,8 +252,12 @@ class BridgedConfigTest {
|
||||
"no profile to launch on ⇒ recognise-only, the pre-CB-557 behaviour");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-579: {@code tab} is the only field a lead's identity depends on now, so it is required
|
||||
* whether the entry is creatable or recognise-only — without it the entry can never be found.
|
||||
*/
|
||||
@Test
|
||||
void aLeadThatCanBeNeitherFoundNorCreatedRefusesToStart(@TempDir Path dir) throws Exception {
|
||||
void aLeadWithNoTabRefusesToStart(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("useless-lead.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
@@ -264,6 +271,80 @@ class BridgedConfigTest {
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, cfg::validateMembers);
|
||||
assertTrue(e.getMessage().contains("ghost"), "the message must name the useless entry");
|
||||
assertTrue(e.getMessage().contains("tab:"), "the message must say what is missing");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-579 acceptance (2): a config still spelling {@code fleet.leaders.<name>.terminal} must fail
|
||||
* loudly at load, not be silently dropped by {@code Leader}'s {@code @JsonIgnoreProperties}.
|
||||
*/
|
||||
@Test
|
||||
void aLeaderTerminalKeyFailsLoadAndNamesTabAsTheReplacement(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("stale-terminal.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
""");
|
||||
|
||||
IllegalStateException e =
|
||||
assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("opus"), "the message must name the offending entry");
|
||||
assertTrue(e.getMessage().contains("tab:"), "the message must name the replacement key");
|
||||
assertTrue(e.getMessage().contains("terminal"), "the message must name the retired key");
|
||||
}
|
||||
|
||||
/** The same refusal, and it must name every offending entry, not just the first. */
|
||||
@Test
|
||||
void everyLeaderStillUsingTerminalIsReportedAtOnce(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("stale-terminals.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
sol:
|
||||
terminal: term_sol
|
||||
""");
|
||||
|
||||
IllegalStateException e =
|
||||
assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("opus"));
|
||||
assertTrue(e.getMessage().contains("sol"));
|
||||
}
|
||||
|
||||
/** A {@code terminal:} anywhere else in the document (not under a leader entry) is unaffected. */
|
||||
@Test
|
||||
void aTerminalKeyOutsideFleetLeadersIsNotRejected(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("primary-terminal-ok.yaml");
|
||||
Files.writeString(f, "bind:\n port: 8080\nprimary:\n terminal: term_fixed\n");
|
||||
|
||||
assertDoesNotThrow(() -> BridgedConfig.load(f));
|
||||
}
|
||||
|
||||
/** CB-579 acceptance (3): distinct `tab:` labels need no shared prefix — one scanner finds both. */
|
||||
@Test
|
||||
void twoLeadersWithDifferentTabsAreBothConfigured(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("two-tabs.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
tab: "lead: opus"
|
||||
sol:
|
||||
tab: "captain: sol"
|
||||
""");
|
||||
|
||||
var leaders = BridgedConfig.load(f).fleet().leaders();
|
||||
assertEquals("lead: opus", leaders.get("opus").tab());
|
||||
assertEquals("captain: sol", leaders.get("sol").tab());
|
||||
}
|
||||
|
||||
// ── CB-551: the idle-lead heartbeat ─────────────────────────────────────────────────────────
|
||||
@@ -328,7 +409,7 @@ class BridgedConfigTest {
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
tab: "lead: opus"
|
||||
tabPrefix: "lead:"
|
||||
""");
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
@@ -349,7 +430,7 @@ class BridgedConfigTest {
|
||||
tabLabel: "lead: {role} {profile}"
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
tab: "lead: opus"
|
||||
""");
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
|
||||
@@ -374,7 +455,7 @@ class BridgedConfigTest {
|
||||
fleet:
|
||||
leaders:
|
||||
opus:
|
||||
terminal: term_opus
|
||||
tab: "lead: opus"
|
||||
""");
|
||||
|
||||
assertDoesNotThrow(() -> BridgedConfig.load(f).validateLeadTabPrefixes());
|
||||
@@ -400,7 +481,7 @@ class BridgedConfigTest {
|
||||
"a label that collides with a convention nobody reads is not a problem");
|
||||
}
|
||||
|
||||
// ── CB-530: the leaders registry ────────────────────────────────────────────────────────────
|
||||
// ── CB-530/CB-579: the leaders registry ─────────────────────────────────────────────────────
|
||||
|
||||
@Test
|
||||
void leadersBlockRegistersEveryPaneByName(@TempDir Path dir) throws Exception {
|
||||
@@ -411,10 +492,10 @@ class BridgedConfigTest {
|
||||
fleet:
|
||||
leaders:
|
||||
opus-5.0:
|
||||
terminal: term_opus
|
||||
tab: "lead: opus-5.0"
|
||||
kind: claude
|
||||
gpt-sol-5.6:
|
||||
terminal: term_sol
|
||||
tab: "lead: gpt-sol-5.6"
|
||||
kind: opencode
|
||||
model: openai/gpt-5.6-terra
|
||||
""");
|
||||
@@ -425,9 +506,9 @@ class BridgedConfigTest {
|
||||
assertEquals(Set.of("opus-5.0", "gpt-sol-5.6"), leaders.keySet());
|
||||
assertEquals("opencode", leaders.get("gpt-sol-5.6").kind());
|
||||
assertEquals("openai/gpt-5.6-terra", leaders.get("gpt-sol-5.6").model());
|
||||
// The whole point: BOTH panes resolve as leads, so neither is demoted to worker.
|
||||
assertEquals(Map.of("term_opus", "opus-5.0", "term_sol", "gpt-sol-5.6"),
|
||||
cfg.leaderTerminals());
|
||||
// Identity is the tab now (CB-579) — both entries carry their own, distinct label.
|
||||
assertEquals("lead: opus-5.0", leaders.get("opus-5.0").tab());
|
||||
assertEquals("lead: gpt-sol-5.6", leaders.get("gpt-sol-5.6").tab());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -439,42 +520,6 @@ class BridgedConfigTest {
|
||||
"configs that never migrate must behave exactly as they did before CB-530");
|
||||
}
|
||||
|
||||
@Test
|
||||
void anExplicitLeadersEntryWinsOverThePinForTheSameTerminal(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("both.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
primary:
|
||||
terminal: term_shared
|
||||
fleet:
|
||||
leaders:
|
||||
opus-5.0:
|
||||
terminal: term_shared
|
||||
""");
|
||||
|
||||
assertEquals(Map.of("term_shared", "opus-5.0"), BridgedConfig.load(f).leaderTerminals(),
|
||||
"the pin is the older spelling of the same fact; the named entry is what was meant");
|
||||
}
|
||||
|
||||
@Test
|
||||
void bothBlocksTogetherRegisterTheUnionOfTheirTerminals(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("union.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
primary:
|
||||
terminal: term_pinned
|
||||
fleet:
|
||||
leaders:
|
||||
gpt-sol-5.6:
|
||||
terminal: term_sol
|
||||
""");
|
||||
|
||||
assertEquals(Map.of("term_pinned", "primary", "term_sol", "gpt-sol-5.6"),
|
||||
BridgedConfig.load(f).leaderTerminals());
|
||||
}
|
||||
|
||||
@Test
|
||||
void neitherBlockLeavesNothingRegistered(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("none.yaml");
|
||||
@@ -483,22 +528,25 @@ class BridgedConfigTest {
|
||||
assertTrue(BridgedConfig.load(f).leaderTerminals().isEmpty());
|
||||
}
|
||||
|
||||
/** A lead entry with no terminal identifies nothing — it must not register a null key. */
|
||||
/**
|
||||
* CB-579: {@code fleet.leaders} no longer feeds {@code leaderTerminals()} at all — a lead's
|
||||
* identity comes from the live tab scan, not a config-held terminal map. This method now exists
|
||||
* only for the {@code primary.terminal} fallback.
|
||||
*/
|
||||
@Test
|
||||
void aLeadWithoutATerminalIsNotRegistered(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-terminal.yaml");
|
||||
void fleetLeadersNeverContributesToLeaderTerminals(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("leaders-only.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
fleet:
|
||||
leaders:
|
||||
sketch:
|
||||
kind: opencode
|
||||
real:
|
||||
terminal: term_real
|
||||
opus-5.0:
|
||||
tab: "lead: opus-5.0"
|
||||
""");
|
||||
|
||||
assertEquals(Map.of("term_real", "real"), BridgedConfig.load(f).leaderTerminals());
|
||||
assertTrue(BridgedConfig.load(f).leaderTerminals().isEmpty(),
|
||||
"no primary.terminal pin ⇒ nothing registered, even with fleet.leaders configured");
|
||||
}
|
||||
|
||||
// ── CB-548: the architects registry ────────────────────────────────────────────────────────
|
||||
|
||||
@@ -0,0 +1,170 @@
|
||||
package dev.ltms.bridged.health;
|
||||
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.bridged.herdr.AgentControl;
|
||||
import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.inject.Injector;
|
||||
import dev.ltms.bridged.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.bridged.msg.MessageService;
|
||||
import dev.ltms.bridged.msg.Rendezvous;
|
||||
import dev.ltms.bridged.session.SessionManager;
|
||||
import dev.ltms.bridged.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||
import dev.ltms.bridged.guard.SubscriptionGuard;
|
||||
import dev.ltms.bridged.config.BridgedConfig;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.function.BiConsumer;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
class FleetHealthMonitorTest {
|
||||
@Test void oneTickUsesOneFleetListForAnyRosterSize() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
BridgedConfig.Profile profile = new BridgedConfig.Profile("test", "http://test:1", null,
|
||||
null, null, null, null, null, null, null, null, null);
|
||||
ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(agents, new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("test")), Map.of("test", profile), "test", _ -> "token");
|
||||
SessionManager sessions = new SessionManager(launcher);
|
||||
sessions.acquire("test", null, null, null);
|
||||
sessions.acquire("test", null, null, null);
|
||||
herdr.calls.clear();
|
||||
MessageService messages = new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox());
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, sessions::roster, messages, scheduler, () -> 1, 60,
|
||||
(_, _) -> { });
|
||||
monitor.tick();
|
||||
monitor.stop();
|
||||
assertEquals(1, herdr.calls.stream().filter(call -> call.method().equals("agent.list")).count());
|
||||
}
|
||||
|
||||
@Test void failedTickDoesNotStopTheNextTick() {
|
||||
FakeHerdr herdr = new FakeHerdr().healthy(false);
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
|
||||
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
|
||||
scheduler, () -> 1, 60, (_, _) -> { });
|
||||
monitor.tick();
|
||||
herdr.healthy(true);
|
||||
monitor.tick();
|
||||
monitor.stop();
|
||||
assertEquals(2, herdr.calls.stream().filter(call -> call.method().equals("agent.list")).count());
|
||||
}
|
||||
|
||||
@Test void faultTransitionLogsOnlyOnceUntilItChanges() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
|
||||
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
|
||||
scheduler, () -> 1, 60, (_, _) -> { });
|
||||
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
|
||||
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
|
||||
monitor.stop();
|
||||
assertEquals(1, appender.list.stream().filter(event -> event.getFormattedMessage()
|
||||
.contains("member=term_a state=TURN_BOUNDARY_LOST")).count());
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
// --- CB-580: a member that reaches GONE/NEVER_READY must fail its waiting tickets
|
||||
|
||||
private static FleetHealthMonitor monitorWith(BiConsumer<String, String> failTarget) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
var scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||
return new FleetHealthMonitor(agents, java.util.List::of,
|
||||
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
|
||||
scheduler, () -> 1, 60, failTarget);
|
||||
}
|
||||
|
||||
@Test void terminalTransitionFailsTheTargetOnce() {
|
||||
RecordingFailTarget failTarget = new RecordingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.stop();
|
||||
assertEquals(1, failTarget.calls.size());
|
||||
assertEquals("term_a", failTarget.calls.get(0).target());
|
||||
assertTrue(failTarget.calls.get(0).reason().contains("GONE"));
|
||||
}
|
||||
|
||||
@Test void neverReadyNamesItselfAsTheReason() {
|
||||
RecordingFailTarget failTarget = new RecordingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
monitor.reportTransition("term_a", HealthState.NEVER_READY);
|
||||
monitor.stop();
|
||||
assertEquals(1, failTarget.calls.size());
|
||||
assertTrue(failTarget.calls.get(0).reason().contains("NEVER_READY"));
|
||||
}
|
||||
|
||||
@Test void stayingInATerminalStateProducesOneFailureNotN() {
|
||||
RecordingFailTarget failTarget = new RecordingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.stop();
|
||||
assertEquals(1, failTarget.calls.size());
|
||||
}
|
||||
|
||||
@Test void aNonTerminalFaultStateDoesNotFailTheTarget() {
|
||||
RecordingFailTarget failTarget = new RecordingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
|
||||
monitor.stop();
|
||||
assertEquals(0, failTarget.calls.size());
|
||||
}
|
||||
|
||||
@Test void failTargetRetryIsBounded() {
|
||||
AlwaysThrowingFailTarget failTarget = new AlwaysThrowingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.stop();
|
||||
assertEquals(FleetHealthMonitor.MAX_FAIL_TARGET_ATTEMPTS, failTarget.calls);
|
||||
}
|
||||
|
||||
@Test void exhaustedRetryStillDoesNotRefireOnAnUnchangedTick() {
|
||||
AlwaysThrowingFailTarget failTarget = new AlwaysThrowingFailTarget();
|
||||
FleetHealthMonitor monitor = monitorWith(failTarget);
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
int afterFirstTransition = failTarget.calls;
|
||||
monitor.reportTransition("term_a", HealthState.GONE);
|
||||
monitor.stop();
|
||||
assertEquals(afterFirstTransition, failTarget.calls);
|
||||
}
|
||||
|
||||
private record RecordedCall(String target, String reason) { }
|
||||
|
||||
private static final class RecordingFailTarget implements BiConsumer<String, String> {
|
||||
final java.util.List<RecordedCall> calls = new java.util.ArrayList<>();
|
||||
|
||||
@Override public void accept(String target, String reason) {
|
||||
calls.add(new RecordedCall(target, reason));
|
||||
}
|
||||
}
|
||||
|
||||
private static final class AlwaysThrowingFailTarget implements BiConsumer<String, String> {
|
||||
int calls = 0;
|
||||
|
||||
@Override public void accept(String target, String reason) {
|
||||
calls++;
|
||||
throw new RuntimeException("boom");
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -15,8 +15,9 @@ import java.util.concurrent.atomic.AtomicLong;
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/**
|
||||
* CB-531. A lead is never spawned, so the daemon has to <em>find</em> it: these assert that an
|
||||
* operator-labelled tab is what makes a pane a lead, and — just as importantly — what does not.
|
||||
* CB-531/CB-579. A lead is never spawned, so the daemon has to <em>find</em> it: these assert that
|
||||
* an operator-labelled tab matching a configured {@code tab:} is what makes a pane a lead, and —
|
||||
* just as importantly — what does not, and that a stale entry does not linger forever.
|
||||
*/
|
||||
class LeadTabScannerTest {
|
||||
|
||||
@@ -120,23 +121,48 @@ class LeadTabScannerTest {
|
||||
.pane("w9:p1", "w9:t1", "term_worker");
|
||||
}
|
||||
|
||||
private LeadTabScanner scanner(TopologyHerdr herdr, Map<String, String> configured,
|
||||
/** The {@code tab:} → name map {@code twoLeads()}'s two lead tabs are configured under. */
|
||||
private static Map<String, String> twoLeadsConfigured() {
|
||||
return Map.of("lead: opus-5.0", "opus-5.0", "lead: gpt-sol-5.6", "gpt-sol-5.6");
|
||||
}
|
||||
|
||||
private LeadTabScanner scanner(TopologyHerdr herdr, Map<String, String> tabToName,
|
||||
AtomicLong clock) {
|
||||
return new LeadTabScanner(herdr, "lead:", Set.of("bridged-workers"), configured, TTL,
|
||||
clock::get);
|
||||
return new LeadTabScanner(herdr, tabToName, Set.of("bridged-workers"), TTL, clock::get);
|
||||
}
|
||||
|
||||
@Test
|
||||
void everyLabelledTabBecomesALeadNamedByItsLabel() {
|
||||
Map<String, String> leads = scanner(twoLeads(), Map.of(), new AtomicLong()).get();
|
||||
void everyConfiguredTabBecomesALeadNamedByItsEntry() {
|
||||
Map<String, String> leads = scanner(twoLeads(), twoLeadsConfigured(), new AtomicLong()).get();
|
||||
|
||||
assertEquals(Map.of("term_opus", "opus-5.0", "term_gpt", "gpt-sol-5.6"), leads,
|
||||
"two leads discovered from labels alone — no terminal_id was ever configured");
|
||||
"two leads discovered by their configured tab — no terminal_id was ever configured");
|
||||
}
|
||||
|
||||
@Test
|
||||
void anUnlabelledTabContributesNothing() {
|
||||
assertFalse(scanner(twoLeads(), Map.of(), new AtomicLong()).get().containsKey("term_notes"));
|
||||
void anUnconfiguredTabContributesNothing() {
|
||||
assertFalse(scanner(twoLeads(), twoLeadsConfigured(), new AtomicLong())
|
||||
.get().containsKey("term_notes"));
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-579: matching is exact against the configured map now, not a shared prefix — two leads with
|
||||
* completely different labels are both discovered by one scanner, no convention required.
|
||||
*/
|
||||
@Test
|
||||
void twoLeadsWithCompletelyDifferentLabelsAreBothDiscovered() {
|
||||
TopologyHerdr herdr = new TopologyHerdr()
|
||||
.workspace("w1", "main")
|
||||
.tab("w1:t1", "w1", "orchestrator: opus")
|
||||
.tab("w1:t2", "w1", "captain: sol")
|
||||
.pane("w1:p1", "w1:t1", "term_opus")
|
||||
.pane("w1:p2", "w1:t2", "term_sol");
|
||||
Map<String, String> tabToName = Map.of("orchestrator: opus", "opus", "captain: sol", "sol");
|
||||
|
||||
Map<String, String> leads = scanner(herdr, tabToName, new AtomicLong()).get();
|
||||
|
||||
assertEquals(Map.of("term_opus", "opus", "term_sol", "sol"), leads,
|
||||
"no shared prefix needed — each lead is matched by its own configured tab");
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -148,25 +174,28 @@ class LeadTabScannerTest {
|
||||
void aTabInAWorkerSpaceIsNeverALeadEvenWhenItsLabelMatches() {
|
||||
TopologyHerdr herdr = twoLeads().tab("w9:t2", "w9", "lead: impostor")
|
||||
.pane("w9:p2", "w9:t2", "term_impostor");
|
||||
Map<String, String> tabToName = new LinkedHashMap<>(twoLeadsConfigured());
|
||||
tabToName.put("lead: impostor", "impostor");
|
||||
|
||||
assertFalse(scanner(herdr, Map.of(), new AtomicLong()).get().containsKey("term_impostor"));
|
||||
assertFalse(scanner(herdr, tabToName, new AtomicLong()).get().containsKey("term_impostor"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aBarePrefixNamesNobodyAndIsRejected() {
|
||||
void aLabelWithNoConfiguredEntryIsIgnored() {
|
||||
TopologyHerdr herdr = new TopologyHerdr().workspace("w1", "main")
|
||||
.tab("w1:t1", "w1", "lead:").pane("w1:p1", "w1:t1", "term_a");
|
||||
.tab("w1:t1", "w1", "lead: nobody-configured").pane("w1:p1", "w1:t1", "term_a");
|
||||
|
||||
assertEquals(Map.of(), scanner(herdr, Map.of(), new AtomicLong()).get(),
|
||||
"a lead with no name would resolve as PRIMARY with nothing to attribute it to");
|
||||
assertEquals(Map.of(), scanner(herdr, twoLeadsConfigured(), new AtomicLong()).get(),
|
||||
"a label that names no configured lead resolves nobody");
|
||||
}
|
||||
|
||||
@Test
|
||||
void thePrefixMatchesCaseInsensitivelyAndTheNameIsTrimmed() {
|
||||
void matchingIsCaseInsensitiveAndToleratesSurroundingWhitespace() {
|
||||
TopologyHerdr herdr = new TopologyHerdr().workspace("w1", "main")
|
||||
.tab("w1:t1", "w1", " LEAD: opus-5.0 ").pane("w1:p1", "w1:t1", "term_a");
|
||||
.tab("w1:t1", "w1", " LEAD: Opus-5.0 ").pane("w1:p1", "w1:t1", "term_a");
|
||||
|
||||
assertEquals(Map.of("term_a", "opus-5.0"), scanner(herdr, Map.of(), new AtomicLong()).get());
|
||||
assertEquals(Map.of("term_a", "opus-5.0"),
|
||||
scanner(herdr, Map.of("lead: Opus-5.0", "opus-5.0"), new AtomicLong()).get());
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -175,18 +204,51 @@ class LeadTabScannerTest {
|
||||
// nothing bridged placed can land here (see the worker-space test above).
|
||||
TopologyHerdr herdr = twoLeads().pane("w1:p1b", "w1:t1", "term_opus_split");
|
||||
|
||||
assertEquals("opus-5.0", scanner(herdr, Map.of(), new AtomicLong()).get().get("term_opus_split"));
|
||||
assertEquals("opus-5.0",
|
||||
scanner(herdr, twoLeadsConfigured(), new AtomicLong()).get().get("term_opus_split"));
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-579 acceptance (6): this is the bug the ticket closes. A stale pin used to be merged back
|
||||
* over every scan and never expire; now a scan is the whole answer, so a lead whose tab is gone
|
||||
* drops out on the very next scan.
|
||||
*/
|
||||
@Test
|
||||
void anExplicitlyConfiguredLeadIsMergedInAndOutranksALabel() {
|
||||
Map<String, String> configured = Map.of("term_opus", "pinned-name", "term_extra", "from-config");
|
||||
void aTabNoLongerPresentDropsTheLeadOnTheNextScan() {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
AtomicLong clock = new AtomicLong();
|
||||
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
|
||||
assertTrue(s.get().containsKey("term_opus"));
|
||||
|
||||
Map<String, String> leads = scanner(twoLeads(), configured, new AtomicLong()).get();
|
||||
// The session behind term_opus restarted — herdr no longer reports that tab or pane at all.
|
||||
herdr.tabs.remove("w1:t1");
|
||||
herdr.panes.remove("w1:p1");
|
||||
clock.addAndGet(TTL);
|
||||
|
||||
assertEquals("pinned-name", leads.get("term_opus"), "an explicit pin is the operator's last word");
|
||||
assertEquals("from-config", leads.get("term_extra"), "a configured lead needs no tab at all");
|
||||
assertEquals("gpt-sol-5.6", leads.get("term_gpt"));
|
||||
assertFalse(s.get().containsKey("term_opus"),
|
||||
"a stale entry must expire once the tab it named is gone, not be merged back forever");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-579 acceptance (5): the whole point of matching by tab instead of {@code terminal_id} — a
|
||||
* restart changes the terminal, not the tab, so the lead resolves under the same name with no
|
||||
* config edit.
|
||||
*/
|
||||
@Test
|
||||
void aLeadRestartingInTheSameTabResolvesUnderTheSameName() {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
AtomicLong clock = new AtomicLong();
|
||||
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
|
||||
assertEquals("opus-5.0", s.get().get("term_opus"));
|
||||
|
||||
// The session restarts: herdr assigns the pane a new terminal_id, same tab (w1:t1).
|
||||
herdr.panes.remove("w1:p1");
|
||||
herdr.pane("w1:p1", "w1:t1", "term_opus_v2");
|
||||
clock.addAndGet(TTL);
|
||||
|
||||
Map<String, String> leads = s.get();
|
||||
assertEquals("opus-5.0", leads.get("term_opus_v2"), "the new terminal resolves immediately");
|
||||
assertFalse(leads.containsKey("term_opus"), "the old terminal_id is simply gone, not carried");
|
||||
}
|
||||
|
||||
// ── caching ─────────────────────────────────────────────────────────────────────────────────
|
||||
@@ -195,7 +257,7 @@ class LeadTabScannerTest {
|
||||
void aSecondLookupWithinTheTtlDoesNotTouchHerdr() {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
AtomicLong clock = new AtomicLong();
|
||||
LeadTabScanner s = scanner(herdr, Map.of(), clock);
|
||||
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
|
||||
|
||||
s.get();
|
||||
int afterFirst = herdr.calls;
|
||||
@@ -210,21 +272,23 @@ class LeadTabScannerTest {
|
||||
void aTabLabelledAfterStartupIsPickedUpOnceTheTtlExpires() {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
AtomicLong clock = new AtomicLong();
|
||||
LeadTabScanner s = scanner(herdr, Map.of(), clock);
|
||||
Map<String, String> tabToName = new LinkedHashMap<>(twoLeadsConfigured());
|
||||
tabToName.put("lead: late-arrival", "late-arrival");
|
||||
LeadTabScanner s = scanner(herdr, tabToName, clock);
|
||||
assertFalse(s.get().containsKey("term_notes"));
|
||||
|
||||
herdr.tab("w1:t3", "w1", "lead: late-arrival"); // the operator renames their tab
|
||||
clock.addAndGet(TTL);
|
||||
|
||||
assertEquals("late-arrival", s.get().get("term_notes"),
|
||||
"the whole point over `leaders:`: no config edit, no restart");
|
||||
"the whole point over a config-held terminal_id: no config edit, no restart");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFailedScanKeepsTheLeadsAlreadyKnownRatherThanDemotingThem() {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
AtomicLong clock = new AtomicLong();
|
||||
LeadTabScanner s = scanner(herdr, Map.of(), clock);
|
||||
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
|
||||
Map<String, String> before = s.get();
|
||||
|
||||
herdr.failing = true;
|
||||
@@ -235,14 +299,14 @@ class LeadTabScannerTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFailedFirstScanStillHonoursTheConfiguredLeads() {
|
||||
void aFailedFirstScanReturnsEmptyRatherThanThrowing() {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
herdr.failing = true;
|
||||
|
||||
Map<String, String> leads = scanner(herdr, Map.of("term_x", "opus-5.0"), new AtomicLong()).get();
|
||||
Map<String, String> leads = scanner(herdr, twoLeadsConfigured(), new AtomicLong()).get();
|
||||
|
||||
assertEquals(Map.of("term_x", "opus-5.0"), leads,
|
||||
"config-named leads must not depend on herdr answering at all");
|
||||
assertEquals(Map.of(), leads,
|
||||
"with nothing scanned yet and no override to fall back on, the map is simply empty");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -250,7 +314,7 @@ class LeadTabScannerTest {
|
||||
TopologyHerdr herdr = twoLeads();
|
||||
herdr.failing = true;
|
||||
AtomicLong clock = new AtomicLong();
|
||||
LeadTabScanner s = scanner(herdr, Map.of(), clock);
|
||||
LeadTabScanner s = scanner(herdr, twoLeadsConfigured(), clock);
|
||||
|
||||
s.get();
|
||||
int afterFirst = herdr.calls;
|
||||
|
||||
@@ -7,6 +7,8 @@ import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.bridged.herdr.AgentControl;
|
||||
import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.msg.Rendezvous;
|
||||
import dev.ltms.bridged.msg.TestTurnTokens;
|
||||
import dev.ltms.bridged.msg.TurnToken;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -47,7 +49,7 @@ class CompletionResolverTest {
|
||||
Rendezvous rendezvous = new Rendezvous(); // no waiter opened
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
|
||||
resolver.captureBaseline("term_a"); // no send to attribute a later completion to
|
||||
resolver.captureBaseline("term_a", TestTurnTokens.inert("term_a")); // no send to attribute a later completion to
|
||||
|
||||
assertFalse(herdr.called("agent.read"),
|
||||
"with no waiting send there is no turn to baseline — skip the scrape");
|
||||
@@ -194,7 +196,7 @@ class CompletionResolverTest {
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.captureBaseline("term_a");
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter));
|
||||
herdr.readText("⏺ answer that /clear would erase\n❯ ");
|
||||
|
||||
resolver.resolveBeforePostAction("term_a");
|
||||
@@ -217,7 +219,7 @@ class CompletionResolverTest {
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
|
||||
|
||||
var waiter = rendezvous.open("term_a"); // a send is blocked on this turn
|
||||
resolver.captureBaseline("term_a"); // baseline is the clipped >cap block
|
||||
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); // baseline is the clipped >cap block
|
||||
var turn = resolver.inFlight("term_a");
|
||||
assertEquals(CompletionResolver.MAX_SCRAPE_CHARS, turn.baseline().length(),
|
||||
"the delivery baseline is clipped to the same cap resolve() applies to the tail");
|
||||
|
||||
@@ -8,6 +8,7 @@ import dev.ltms.bridged.herdr.AgentControl;
|
||||
import dev.ltms.bridged.herdr.AgentStatus;
|
||||
import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.herdr.HerdrException;
|
||||
import dev.ltms.bridged.msg.TestTurnTokens;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -47,7 +48,7 @@ class InjectorTest {
|
||||
|
||||
@Test
|
||||
void deliversWhenIdle() {
|
||||
CompletableFuture<Void> f = injector.enqueue(T, "hello");
|
||||
CompletableFuture<Void> f = injector.enqueue(T, "hello", TestTurnTokens.inert(T));
|
||||
assertFalse(f.isDone(), "not delivered until an injectable status arrives");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
assertTrue(f.isDone());
|
||||
@@ -59,7 +60,7 @@ class InjectorTest {
|
||||
// CB-113: idle alone is not enough — hold until the worker's MCP is connected (ready).
|
||||
java.util.Set<String> ready = new java.util.HashSet<>();
|
||||
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, ready::contains);
|
||||
inj.enqueue(T, "task");
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // idle but not yet available → held out of the boot window
|
||||
assertEquals(List.of(), sent(), "must not deliver into a not-yet-available worker");
|
||||
@@ -79,7 +80,7 @@ class InjectorTest {
|
||||
void resubmitsEnterWhenADeliveredMessageIsNotPickedUp() {
|
||||
// CB-113: the Enter at delivery can race the paste; while the worker stays idle (not picked
|
||||
// up), the injector re-nudges Enter so the pending paste submits.
|
||||
injector.enqueue(T, "task");
|
||||
injector.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver: paste + one Enter
|
||||
long afterDeliver = enterKeystrokes();
|
||||
|
||||
@@ -95,7 +96,7 @@ class InjectorTest {
|
||||
|
||||
@Test
|
||||
void holdsWhileWorkingThenDeliversOnIdle() {
|
||||
injector.enqueue(T, "later");
|
||||
injector.enqueue(T, "later", TestTurnTokens.inert(T));
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
assertEquals(List.of(), sent(), "must not inject mid-turn");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
@@ -104,7 +105,7 @@ class InjectorTest {
|
||||
|
||||
@Test
|
||||
void blockedIsInjectableButUnknownIsNot() {
|
||||
injector.enqueue(T, "answer");
|
||||
injector.enqueue(T, "answer", TestTurnTokens.inert(T));
|
||||
injector.onStatus(T, AgentStatus.UNKNOWN);
|
||||
assertEquals(List.of(), sent(), "unknown status is not safe to inject");
|
||||
injector.onStatus(T, AgentStatus.BLOCKED);
|
||||
@@ -113,8 +114,8 @@ class InjectorTest {
|
||||
|
||||
@Test
|
||||
void twoRapidDeliveriesNeverInterleave() {
|
||||
injector.enqueue(T, "m1");
|
||||
injector.enqueue(T, "m2");
|
||||
injector.enqueue(T, "m1", TestTurnTokens.inert(T));
|
||||
injector.enqueue(T, "m2", TestTurnTokens.inert(T));
|
||||
|
||||
// First idle window delivers only m1, even if idle is observed twice before pickup.
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
@@ -129,8 +130,8 @@ class InjectorTest {
|
||||
|
||||
@Test
|
||||
void transientUnknownDoesNotReleaseThePickupLatch() {
|
||||
injector.enqueue(T, "m1");
|
||||
injector.enqueue(T, "m2");
|
||||
injector.enqueue(T, "m1", TestTurnTokens.inert(T));
|
||||
injector.enqueue(T, "m2", TestTurnTokens.inert(T));
|
||||
injector.onStatus(T, AgentStatus.IDLE); // m1 sent, awaiting pickup
|
||||
assertEquals(List.of("m1"), sent());
|
||||
|
||||
@@ -145,8 +146,8 @@ class InjectorTest {
|
||||
|
||||
@Test
|
||||
void missedPickupEdgeIsReleasedByGraceSoTheQueueNeverWedges() {
|
||||
injector.enqueue(T, "m1");
|
||||
injector.enqueue(T, "m2");
|
||||
injector.enqueue(T, "m1", TestTurnTokens.inert(T));
|
||||
injector.enqueue(T, "m2", TestTurnTokens.inert(T));
|
||||
injector.onStatus(T, AgentStatus.IDLE); // m1 sent
|
||||
assertEquals(List.of("m1"), sent());
|
||||
|
||||
@@ -158,9 +159,9 @@ class InjectorTest {
|
||||
|
||||
@Test
|
||||
void fifoOrderAcrossManyTurns() {
|
||||
injector.enqueue(T, "a");
|
||||
injector.enqueue(T, "b");
|
||||
injector.enqueue(T, "c");
|
||||
injector.enqueue(T, "a", TestTurnTokens.inert(T));
|
||||
injector.enqueue(T, "b", TestTurnTokens.inert(T));
|
||||
injector.enqueue(T, "c", TestTurnTokens.inert(T));
|
||||
for (int i = 0; i < 3; i++) {
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver one
|
||||
injector.onStatus(T, AgentStatus.WORKING); // pickup
|
||||
@@ -172,7 +173,7 @@ class InjectorTest {
|
||||
@Test
|
||||
void activeWhileQueuedOrInFlightThenQuietAfterTurnCompletes() {
|
||||
assertTrue(injector.activeTargets().isEmpty());
|
||||
injector.enqueue(T, "x");
|
||||
injector.enqueue(T, "x", TestTurnTokens.inert(T));
|
||||
assertEquals(Set.of(T), injector.activeTargets(), "active while a message is queued");
|
||||
|
||||
injector.onStatus(T, AgentStatus.IDLE); // delivers; awaiting pickup
|
||||
@@ -191,7 +192,7 @@ class InjectorTest {
|
||||
void firesTurnCompleteOnAConfirmedWorkingThenIdle() {
|
||||
List<String> completed = new ArrayList<>();
|
||||
Injector inj = new Injector(new AgentControl(herdr), completed::add);
|
||||
inj.enqueue(T, "task");
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
inj.onStatus(T, AgentStatus.WORKING); // pickup + turn running
|
||||
@@ -226,8 +227,8 @@ class InjectorTest {
|
||||
}
|
||||
ResetListener listener = new ResetListener();
|
||||
Injector inj = new Injector(agents, listener);
|
||||
inj.enqueue(T, "first");
|
||||
inj.enqueue(T, "second");
|
||||
inj.enqueue(T, "first", TestTurnTokens.inert(T));
|
||||
inj.enqueue(T, "second", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // first delegation
|
||||
inj.onStatus(T, AgentStatus.WORKING);
|
||||
@@ -246,7 +247,7 @@ class InjectorTest {
|
||||
void doesNotSynthesizeCompletionFromAnUnconfirmedTurn() {
|
||||
List<String> completed = new ArrayList<>();
|
||||
Injector inj = new Injector(new AgentControl(herdr), completed::add);
|
||||
inj.enqueue(T, "task");
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
// Deliver, then only ever idle — a `working` sample is never seen. The pickup grace unwedges
|
||||
// the queue but must NOT invent a completion: without a sampled turn there is no trustworthy
|
||||
@@ -259,6 +260,7 @@ class InjectorTest {
|
||||
private static final class Captor implements TurnListener {
|
||||
final List<String> completed = new ArrayList<>();
|
||||
final List<String> failed = new ArrayList<>();
|
||||
final List<String> failureReasons = new ArrayList<>();
|
||||
|
||||
@Override
|
||||
public void onTurnComplete(String target) {
|
||||
@@ -269,6 +271,12 @@ class InjectorTest {
|
||||
public void onTurnFailed(String target) {
|
||||
failed.add(target);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onTurnFailed(String target, String reason) {
|
||||
failed.add(target);
|
||||
failureReasons.add(reason);
|
||||
}
|
||||
}
|
||||
|
||||
// ~30s of unknown at the 250ms prod poll interval; enough onStatus samples to trip the stall.
|
||||
@@ -278,7 +286,7 @@ class InjectorTest {
|
||||
void failsAnOutstandingDelegationWhoseWorkerWedgesInUnknown() {
|
||||
Captor cap = new Captor();
|
||||
Injector inj = new Injector(new AgentControl(herdr), cap);
|
||||
inj.enqueue(T, "task");
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
inj.onStatus(T, AgentStatus.WORKING); // worker starts the turn
|
||||
@@ -293,7 +301,7 @@ class InjectorTest {
|
||||
void aTransientUnknownGlitchNeitherFailsNorBlocksCompletion() {
|
||||
Captor cap = new Captor();
|
||||
Injector inj = new Injector(new AgentControl(herdr), cap);
|
||||
inj.enqueue(T, "task");
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
inj.onStatus(T, AgentStatus.WORKING); // confirmed turn
|
||||
@@ -308,7 +316,7 @@ class InjectorTest {
|
||||
void sendFailureDropsMessageAndFailsItsFuture() {
|
||||
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
|
||||
Injector inj = new Injector(new AgentControl(failing));
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "boom");
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
assertTrue(f.isCompletedExceptionally());
|
||||
@@ -317,18 +325,35 @@ class InjectorTest {
|
||||
|
||||
@Test
|
||||
void dropFailsPendingWaiters() {
|
||||
CompletableFuture<Void> f = injector.enqueue(T, "orphan");
|
||||
CompletableFuture<Void> f = injector.enqueue(T, "orphan", TestTurnTokens.inert(T));
|
||||
injector.drop(T, new HerdrException("worker gone", "pane_not_found", null));
|
||||
assertTrue(f.isCompletedExceptionally(), "queued waiters unblock when the worker vanishes");
|
||||
}
|
||||
|
||||
@Test
|
||||
void dropPassesTheRealCauseForQueuedAndDeliveredWork() {
|
||||
Captor cap = new Captor();
|
||||
Injector inj = new Injector(new AgentControl(herdr), cap);
|
||||
CompletableFuture<Void> delivered = inj.enqueue(T, "delivered", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> queued = inj.enqueue(T, "queued", TestTurnTokens.inert(T));
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // deliver the first message
|
||||
inj.onStatus(T, AgentStatus.WORKING); // its turn is now in flight; one remains queued
|
||||
inj.drop(T, new HerdrException("agent target sol not found", "agent_not_found", null));
|
||||
|
||||
assertEquals(List.of(T), cap.failed, "drop signals one turn failure for both affected states");
|
||||
assertEquals(List.of("agent target sol not found"), cap.failureReasons);
|
||||
assertTrue(delivered.isDone(), "the delivered future has already completed");
|
||||
assertTrue(queued.isCompletedExceptionally(), "the queued future fails with the drop cause");
|
||||
}
|
||||
|
||||
@Test
|
||||
void dropFailsTheTurnOfADeliveredMessageWhenTheWorkerVanishes() {
|
||||
// CB-110: the message was delivered (no longer queued), so failing queued waiters alone would
|
||||
// leave its send hanging. A vanished worker must fail that in-flight turn too.
|
||||
Captor cap = new Captor();
|
||||
Injector inj = new Injector(new AgentControl(herdr), cap);
|
||||
inj.enqueue(T, "task");
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
inj.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
inj.onStatus(T, AgentStatus.WORKING); // turn running
|
||||
|
||||
@@ -343,7 +368,7 @@ class InjectorTest {
|
||||
// by an off-sub worker's review of CB-110, delegated through the bridge.)
|
||||
Captor cap = new Captor();
|
||||
Injector inj = new Injector(new AgentControl(herdr), cap);
|
||||
inj.enqueue(T, "task");
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
inj.onStatus(T, AgentStatus.IDLE); // deliver; pickup never confirmed
|
||||
|
||||
inj.drop(T, new HerdrException("worker gone", "pane_not_found", null));
|
||||
@@ -362,7 +387,7 @@ class InjectorTest {
|
||||
Captor cap = new Captor();
|
||||
List<String> forgotten = new ArrayList<>();
|
||||
Injector inj = new Injector(new AgentControl(herdr), cap, _ -> false, forgotten::add);
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "task");
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
@@ -390,7 +415,7 @@ class InjectorTest {
|
||||
try {
|
||||
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> false, _ -> {
|
||||
});
|
||||
inj.enqueue(T, "task");
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
@@ -415,7 +440,7 @@ class InjectorTest {
|
||||
Set<String> ready = new java.util.HashSet<>();
|
||||
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, ready::contains, _ -> {
|
||||
});
|
||||
inj.enqueue(T, "task");
|
||||
inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
|
||||
for (int i = 0; i < 100; i++) inj.onStatus(T, AgentStatus.IDLE); // still booting, well under grace
|
||||
assertEquals(List.of(), sent());
|
||||
@@ -431,7 +456,7 @@ class InjectorTest {
|
||||
// linger past the worker's life (MemberPresence.forget had no caller before this).
|
||||
List<String> forgotten = new ArrayList<>();
|
||||
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> true, forgotten::add);
|
||||
inj.enqueue(T, "orphan");
|
||||
inj.enqueue(T, "orphan", TestTurnTokens.inert(T));
|
||||
inj.drop(T, new HerdrException("worker gone", "pane_not_found", null));
|
||||
assertEquals(List.of(T), forgotten, "drop clears the gone worker's presence");
|
||||
}
|
||||
@@ -452,7 +477,7 @@ class InjectorTest {
|
||||
try {
|
||||
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> true, _ -> {
|
||||
});
|
||||
inj.enqueue(T, "orphan");
|
||||
inj.enqueue(T, "orphan", TestTurnTokens.inert(T));
|
||||
inj.drop(T, new HerdrException("worker gone", "pane_not_found", null));
|
||||
|
||||
String warn = appender.list.stream()
|
||||
@@ -476,7 +501,7 @@ class InjectorTest {
|
||||
StatusPoller poller = new StatusPoller(new AgentControl(idle), inj, 10);
|
||||
poller.start();
|
||||
try {
|
||||
CompletableFuture<Void> delivered = inj.enqueue(T, "via-poller");
|
||||
CompletableFuture<Void> delivered = inj.enqueue(T, "via-poller", TestTurnTokens.inert(T));
|
||||
delivered.get(2, TimeUnit.SECONDS); // completes when the poller drives the send
|
||||
} finally {
|
||||
poller.stop();
|
||||
@@ -495,7 +520,7 @@ class InjectorTest {
|
||||
void deliveredFutureCarriesSendFailure() {
|
||||
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
|
||||
Injector inj = new Injector(new AgentControl(failing));
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "boom");
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
ExecutionException ex = assertThrows(ExecutionException.class, f::get);
|
||||
assertInstanceOf(HerdrException.class, ex.getCause());
|
||||
|
||||
@@ -41,8 +41,8 @@ class LeadLauncherTest {
|
||||
null, null, fleet, null, "fixed", null).withDefaults();
|
||||
}
|
||||
|
||||
private static BridgedConfig.Leader lead(String profile, String terminal, int instances) {
|
||||
return new BridgedConfig.Leader(profile, terminal, instances, "lead:", 10, null, null,
|
||||
private static BridgedConfig.Leader lead(String profile, String tab, int instances) {
|
||||
return new BridgedConfig.Leader(profile, tab, instances, "lead:", 10, null, null,
|
||||
"leads", "/repo");
|
||||
}
|
||||
|
||||
@@ -70,16 +70,16 @@ class LeadLauncherTest {
|
||||
void startsTheDeclaredLeadWhenNoneIsRunning() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads());
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
|
||||
assertTrue(herdr.called("agent.start"), "a lead must actually be started");
|
||||
assertEquals("lead-opus", startedName(herdr));
|
||||
}
|
||||
|
||||
/** The tab is labelled so the scanner finds the lead on the next resolve. */
|
||||
/** The tab is labelled with the configured `tab:` so the scanner finds the lead on the next resolve. */
|
||||
@Test
|
||||
void labelsTheTabWithThePrefixTheScannerReadsBack() {
|
||||
void labelsTheTabWithTheConfiguredTabValue() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads();
|
||||
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads();
|
||||
|
||||
assertEquals("lead: opus",
|
||||
((Map<?, ?>) herdr.lastCall("tab.rename").params()).get("label"));
|
||||
@@ -90,7 +90,7 @@ class LeadLauncherTest {
|
||||
void startsAsManyInstancesAsAreDeclared() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
assertEquals(2, launcher(herdr, configWith(lead("opus", null, 2))).ensureLeads());
|
||||
assertEquals(2, launcher(herdr, configWith(lead("opus", "lead: opus", 2))).ensureLeads());
|
||||
assertEquals(2, herdr.calls.stream().filter(c -> c.method().equals("agent.start")).count());
|
||||
}
|
||||
|
||||
@@ -104,7 +104,7 @@ class LeadLauncherTest {
|
||||
.withTab("wL", "wL:t1", "lead: opus")
|
||||
.withAgent("lead-opus", "term_lead", "wL:p1", "wL:t1");
|
||||
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads());
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
|
||||
assertFalse(herdr.called("agent.start"), "the live lead must not be duplicated");
|
||||
}
|
||||
|
||||
@@ -118,20 +118,23 @@ class LeadLauncherTest {
|
||||
.withWorkspace("wL", "leads")
|
||||
.withTab("wL", "wL:t1", "lead: opus"); // label only — nothing running in it
|
||||
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads(),
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
|
||||
"a stale label is not a lead; the lead must be relaunched");
|
||||
}
|
||||
|
||||
/**
|
||||
* A lead the operator opened by hand and pinned with `terminal:` is live even though its tab
|
||||
* carries no matching label. Counting labels alone would relaunch it on every boot.
|
||||
* A lead the operator opened by hand is live once its tab carries the configured `tab:` label —
|
||||
* CB-579 retired the `terminal:` pin, so a hand-opened lead is found the same way an
|
||||
* auto-launched one is, by its tab, not by a terminal id nobody wrote down in advance.
|
||||
*/
|
||||
@Test
|
||||
void aPinnedTerminalWithARunningAgentCountsAsLive() {
|
||||
void aHandOpenedLeadWithTheConfiguredTabLabelCountsAsLive() {
|
||||
FakeHerdr herdr = new FakeHerdr()
|
||||
.withAgent("hand-opened", "term_pinned", "wX:p1", "wX:t1");
|
||||
.withWorkspace("wX", "main")
|
||||
.withTab("wX", "wX:t1", "lead: opus")
|
||||
.withAgent("hand-opened", "term_hand", "wX:p1", "wX:t1");
|
||||
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", "term_pinned", 1))).ensureLeads());
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
|
||||
assertFalse(herdr.called("agent.start"));
|
||||
}
|
||||
|
||||
@@ -143,7 +146,7 @@ class LeadLauncherTest {
|
||||
.withTab("wM", "wM:t1", "lead: opus") // a member tab that looks like a lead
|
||||
.withAgent("claude-opus-x", "term_m", "wM:p1", "wM:t1");
|
||||
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads(),
|
||||
assertEquals(1, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads(),
|
||||
"a member in a lead-labelled tab is not a lead, so the real lead is still missing");
|
||||
}
|
||||
|
||||
@@ -152,7 +155,7 @@ class LeadLauncherTest {
|
||||
void anUncountableHerdrStartsNothing() {
|
||||
FakeHerdr herdr = new FakeHerdr().healthy(false);
|
||||
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads());
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads());
|
||||
assertFalse(herdr.called("agent.start"));
|
||||
}
|
||||
|
||||
@@ -166,7 +169,7 @@ class LeadLauncherTest {
|
||||
@Test
|
||||
void theLeadNeverReceivesTheWorkerReplyCharter() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads();
|
||||
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads();
|
||||
|
||||
List<String> args = startedArgs(herdr);
|
||||
assertFalse(args.contains("--append-system-prompt"),
|
||||
@@ -178,7 +181,7 @@ class LeadLauncherTest {
|
||||
@Test
|
||||
void theLeadMountsTheBridgeMcpAndPinsItsModel() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads();
|
||||
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads();
|
||||
|
||||
List<String> args = startedArgs(herdr);
|
||||
assertTrue(args.contains("--mcp-config"));
|
||||
@@ -192,7 +195,7 @@ class LeadLauncherTest {
|
||||
@Test
|
||||
void theLeadEnvCarriesNoAnthropicBinding() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads();
|
||||
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads();
|
||||
|
||||
Map<String, String> env = tabEnv(herdr);
|
||||
assertNull(env.get("ANTHROPIC_BASE_URL"));
|
||||
@@ -205,7 +208,7 @@ class LeadLauncherTest {
|
||||
@Test
|
||||
void theLeadTabIsCreatedOutsideEveryMemberWorkspace() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
launcher(herdr, configWith(lead("opus", null, 1))).ensureLeads();
|
||||
launcher(herdr, configWith(lead("opus", "lead: opus", 1))).ensureLeads();
|
||||
|
||||
String label = (String) ((Map<?, ?>) herdr.lastCall("workspace.create").params()).get("label");
|
||||
assertEquals("leads", label);
|
||||
@@ -214,12 +217,12 @@ class LeadLauncherTest {
|
||||
|
||||
// ── recognise-only and misconfiguration ───────────────────────────────────────────────────
|
||||
|
||||
/** A lead with a pin but no profile is recognise-only by design — not an error, not a launch. */
|
||||
/** A lead with a tab but no profile is recognise-only by design — not an error, not a launch. */
|
||||
@Test
|
||||
void aLeadThatNamesNoProfileIsRecognisedButNeverLaunched() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
assertEquals(0, launcher(herdr, configWith(lead(null, "term_dead", 1))).ensureLeads());
|
||||
assertEquals(0, launcher(herdr, configWith(lead(null, "lead: dead", 1))).ensureLeads());
|
||||
assertFalse(herdr.called("agent.start"));
|
||||
}
|
||||
|
||||
@@ -228,7 +231,7 @@ class LeadLauncherTest {
|
||||
void zeroInstancesLaunchesNothing() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", null, 0))).ensureLeads());
|
||||
assertEquals(0, launcher(herdr, configWith(lead("opus", "lead: opus", 0))).ensureLeads());
|
||||
assertFalse(herdr.called("agent.start"));
|
||||
}
|
||||
|
||||
@@ -237,7 +240,7 @@ class LeadLauncherTest {
|
||||
void anUnknownProfileIsSkippedRatherThanThrown() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
assertEquals(0, launcher(herdr, configWith(lead("nope", null, 1))).ensureLeads());
|
||||
assertEquals(0, launcher(herdr, configWith(lead("nope", "lead: opus", 1))).ensureLeads());
|
||||
assertFalse(herdr.called("agent.start"));
|
||||
}
|
||||
|
||||
|
||||
@@ -72,7 +72,7 @@ class BridgeMcpAuthzTest {
|
||||
new PrimaryRegistry(null),
|
||||
enforce ? CallerResolver.withLeadsAndMembers(identity, false, null,
|
||||
Map::of, new MemberRegistry(null)) : null,
|
||||
metrics, BridgeMcp.CapacitySource.none());
|
||||
metrics, BridgeMcp.CapacitySource.none(), new BridgeMcp.HealthCoverageSource(() -> "off"));
|
||||
return mcp;
|
||||
}
|
||||
|
||||
|
||||
@@ -474,7 +474,7 @@ class BridgeMcpTest {
|
||||
sessions.acquire("ltms-local", null, null, null);
|
||||
McpSchema.CallToolResult res = BridgeMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
|
||||
sessions, null, new BridgeMcp.CapacitySource(profile -> 2, profile -> 2,
|
||||
() -> Set.of("ltms-local"), () -> 0), Map.of(), "");
|
||||
() -> Set.of("ltms-local"), () -> 0), new BridgeMcp.HealthCoverageSource(() -> "off"), Map.of(), "");
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"maxLoad\":2"), out);
|
||||
assertTrue(out.contains("\"live\":2"), out);
|
||||
@@ -487,7 +487,7 @@ class BridgeMcpTest {
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
String out = textOf(BridgeMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
|
||||
sessions, null, new BridgeMcp.CapacitySource(profile -> 0, profile -> 2,
|
||||
() -> Set.of("terra"), () -> 0), Map.of(), ""));
|
||||
() -> Set.of("terra"), () -> 0), new BridgeMcp.HealthCoverageSource(() -> "off"), Map.of(), ""));
|
||||
assertTrue(out.contains("\"profile\":\"terra\""), out);
|
||||
assertTrue(out.contains("\"live\":0"), out);
|
||||
assertTrue(out.contains("\"free\":2"), out);
|
||||
@@ -499,7 +499,7 @@ class BridgeMcpTest {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
String out = textOf(BridgeMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
|
||||
new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))), null,
|
||||
BridgeMcp.CapacitySource.none(), Map.of(), ""));
|
||||
BridgeMcp.CapacitySource.none(), new BridgeMcp.HealthCoverageSource(() -> "off"), Map.of(), ""));
|
||||
assertFalse(out.contains("\"capacity\":"), out);
|
||||
}
|
||||
|
||||
|
||||
@@ -120,6 +120,36 @@ class MessageServiceTest {
|
||||
assertFalse(reply.completed());
|
||||
}
|
||||
|
||||
@Test
|
||||
void droppedQueuedAndDeliveredTurnsExposeTheRealCauseExactlyOnce() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> first = sendAsync();
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // first delivery
|
||||
injector.onStatus(T, AgentStatus.WORKING); // first turn in flight
|
||||
|
||||
CompletableFuture<Void> queued = injector.enqueue(T, "second task", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(T);
|
||||
injector.drop(T, new HerdrException("agent target sol not found", "agent_not_found", null));
|
||||
|
||||
MessageService.Reply reply = first.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.WORKER_FAILED, reply.outcome());
|
||||
assertEquals("agent target sol not found", reply.text());
|
||||
assertTrue(queued.isCompletedExceptionally(), "the queued delivery future also fails");
|
||||
assertFalse(rendezvous.resolveFailure(waiter, "second failure"), "the waiter fails exactly once");
|
||||
}
|
||||
|
||||
@Test
|
||||
void noDropReasonKeepsTheExistingFallbackText() throws Exception {
|
||||
herdr.readText("");
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open(T);
|
||||
|
||||
completion.onTurnFailed(T);
|
||||
|
||||
Rendezvous.Resolution resolution = waiter.get(2, TimeUnit.SECONDS);
|
||||
assertEquals("worker did not reply; its turn ended in an unrecoverable state "
|
||||
+ "(worker unreachable or stuck)", resolution.text());
|
||||
}
|
||||
|
||||
// --- bridge_ask reverse rendezvous (CB-205) ------------------------------------------------
|
||||
|
||||
@Test
|
||||
@@ -609,4 +639,94 @@ class MessageServiceTest {
|
||||
assertTrue(view.detail() != null && view.detail().contains("released"),
|
||||
"and the detail says why, rather than 'worker unknown'");
|
||||
}
|
||||
|
||||
@Test
|
||||
void abandonFailsEveryPendingAsyncTicketForTheReleasedTarget() throws Exception {
|
||||
String first = messages.sendAsync(T, "first task");
|
||||
awaitWaiting(); // first task owns the target lock and rendezvous waiter
|
||||
String second = messages.sendAsync(T, "second task"); // parked on the same lock, not yet queued
|
||||
String third = messages.sendAsync(T, "third task"); // a second queued ticket proves the full sweep
|
||||
|
||||
assertTrue(messages.abandon(T, "agent target term_a not found"));
|
||||
|
||||
assertFailedTicket(first, "agent target term_a not found");
|
||||
assertFailedTicket(second, "agent target term_a not found");
|
||||
assertFailedTicket(third, "agent target term_a not found");
|
||||
}
|
||||
|
||||
@Test
|
||||
void abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
assertFalse(messages.abandon(T, "agent target term_a not found"),
|
||||
"an asking ticket is an active turn, not a pending send to sweep");
|
||||
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase());
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
@Test
|
||||
void unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
assertEquals(MessageService.AskOutcome.TIMED_OUT,
|
||||
messages.ask(T, "which config?", 200).outcome());
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(),
|
||||
"only the question wait ended; the delegated turn may still finish");
|
||||
|
||||
String next = messages.sendAsync(T, "next task");
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
assertEquals(MessageService.Phase.DONE, awaitTicketPhase(next, MessageService.Phase.DONE).phase());
|
||||
}
|
||||
|
||||
@Test
|
||||
void asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter() throws Exception {
|
||||
String first = messages.sendAsync(T, "first task");
|
||||
awaitWaiting();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
assertEquals(MessageService.Phase.ASKING, awaitTicketPhase(first, MessageService.Phase.ASKING).phase());
|
||||
|
||||
MessageService.TaskView asking = messages.poll(first);
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
private void assertFailedTicket(String ticket, String reason) throws Exception {
|
||||
MessageService.TaskView view = awaitTicketPhase(ticket, MessageService.Phase.FAILED);
|
||||
assertEquals(reason, view.detail());
|
||||
}
|
||||
|
||||
private MessageService.TaskView awaitTicketPhase(String ticket, MessageService.Phase phase) throws Exception {
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
MessageService.TaskView view;
|
||||
do {
|
||||
view = messages.poll(ticket);
|
||||
if (view.phase() == phase) {
|
||||
return view;
|
||||
}
|
||||
Thread.sleep(5);
|
||||
} while (System.currentTimeMillis() < deadline);
|
||||
assertEquals(phase, view.phase());
|
||||
return view;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,20 @@
|
||||
package dev.ltms.bridged.msg;
|
||||
|
||||
/**
|
||||
* Explicit unbound tokens for tests that exercise delivery without an accepted send.
|
||||
*
|
||||
* <p>The waiter is {@code null} on purpose. "No accepted send" is an <em>absence</em>, and a helper
|
||||
* that handed back a fresh {@code CompletableFuture} would invent one — which is how the first
|
||||
* version of this class turned {@code captureBaselineSkipsTheReadWhenNoSendIsWaiting} red: the
|
||||
* resolver saw a non-null waiter, decided a turn was in flight, and scraped a pane that no send was
|
||||
* blocked on. An inert value must omit the fact, never fabricate it.
|
||||
*/
|
||||
public final class TestTurnTokens {
|
||||
private TestTurnTokens() {
|
||||
}
|
||||
|
||||
/** A token for a delivery that no send is waiting on: it authorises nothing. */
|
||||
public static TurnToken inert(String target) {
|
||||
return new TurnToken(target, null);
|
||||
}
|
||||
}
|
||||
@@ -29,6 +29,7 @@ public final class FakeWorktrees implements Worktrees {
|
||||
private final Set<String> existingPaths = ConcurrentHashMap.newKeySet();
|
||||
private final Set<String> trackedPaths = ConcurrentHashMap.newKeySet();
|
||||
private volatile RuntimeException addFailure;
|
||||
private volatile boolean dirty = false;
|
||||
private volatile String repoRoot = "/repo";
|
||||
private volatile String prefix = "/worktrees";
|
||||
|
||||
@@ -61,6 +62,12 @@ public final class FakeWorktrees implements Worktrees {
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Mark the worktree dirty so {@link #hasUncommitted} reports true (simulates uncommitted work). */
|
||||
public FakeWorktrees withDirty(boolean dirty) {
|
||||
this.dirty = dirty;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String add(String repoRoot, String branch, String baseRef) {
|
||||
addCalls.add(new AddCall(repoRoot, branch, baseRef));
|
||||
@@ -77,6 +84,11 @@ public final class FakeWorktrees implements Worktrees {
|
||||
removeCalls.add(new RemoveCall(repoRoot, worktreePath));
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean hasUncommitted(String worktreePath) {
|
||||
return dirty;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
|
||||
List<String> copied = new java.util.ArrayList<>();
|
||||
|
||||
@@ -175,6 +175,46 @@ class GitWorktreesTest {
|
||||
assertTrue(Files.exists(Path.of(wt).resolve(".mcp.json")), ".mcp.json stub was dropped");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-576. {@code hasUncommitted} must treat a freshly-provisioned worktree as clean, but a
|
||||
* worktree holding a brand-new, never-added file as dirty. The untracked-file-only shape is
|
||||
* exactly the work lost in the incident — a worker's draft that compiled but was never
|
||||
* committed because it stopped to ask its lead a question.
|
||||
*/
|
||||
@Test
|
||||
void anUntrackedOnlyWorktreeCountsAsDirty(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo"));
|
||||
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
|
||||
String wt = gitWorktrees.add(repo.toString(), "cb-576-u", "HEAD");
|
||||
|
||||
assertFalse(gitWorktrees.hasUncommitted(wt),
|
||||
"a freshly provisioned worktree must read as clean");
|
||||
|
||||
Files.writeString(Path.of(wt).resolve("brand-new.txt"), "draft that was never added\n");
|
||||
|
||||
assertTrue(gitWorktrees.hasUncommitted(wt),
|
||||
"an untracked-only file must count as dirty");
|
||||
|
||||
Files.writeString(Path.of(wt).resolve("README.md"), "edited tracked file\n");
|
||||
assertTrue(gitWorktrees.hasUncommitted(wt),
|
||||
"a tracked modification must also count as dirty");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-576 review. {@code hasUncommitted} must tolerate a missing worktree exactly like
|
||||
* {@code remove}: an already-gone directory holds no work to lose, and throwing here would
|
||||
* break teardown — SessionManager.release() calls it before stopping the pane, so an
|
||||
* exception would orphan a live pane and skip the release notification (CB-516).
|
||||
*/
|
||||
@Test
|
||||
void hasUncommittedOnAMissingWorktreeReturnsFalseWithoutThrowing(@TempDir Path tmp) {
|
||||
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
|
||||
String gone = tmp.resolve("wts").resolve("does-not-exist").toString();
|
||||
|
||||
assertFalse(gitWorktrees.hasUncommitted(gone),
|
||||
"a missing worktree is reported clean, not an error");
|
||||
}
|
||||
|
||||
/** All three protected configs are covered: each one present in a worktree is neutralized and hidden. */
|
||||
@Test
|
||||
void allThreeConfigsAreNeutralizedWhenPresent(@TempDir Path tmp) throws Exception {
|
||||
|
||||
@@ -10,6 +10,7 @@ import dev.ltms.bridged.herdr.AgentControl;
|
||||
import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||
import dev.ltms.bridged.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.bridged.msg.TestTurnTokens;
|
||||
import dev.ltms.bridged.peer.CharterReceipt;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
import dev.ltms.bridged.peer.PeerUnreachableException;
|
||||
@@ -121,7 +122,7 @@ class SessionManagerTest {
|
||||
|
||||
assertDoesNotThrow(() -> sessions.asPresence().markPresent(null),
|
||||
"the primary's null terminal must not blow up an unrelated tool call");
|
||||
assertDoesNotThrow(() -> sessions.onDelivered(null));
|
||||
assertDoesNotThrow(() -> sessions.onDelivered(null, TestTurnTokens.inert(null)));
|
||||
assertDoesNotThrow(() -> sessions.onTurnComplete(null));
|
||||
assertDoesNotThrow(() -> sessions.onTurnFailed(null));
|
||||
|
||||
@@ -141,7 +142,7 @@ class SessionManagerTest {
|
||||
"MCP presence moves SPAWNING → READY");
|
||||
assertTrue(sessions.asPresence().isPresent(terminal), "presence is also recorded");
|
||||
|
||||
sessions.onDelivered(terminal);
|
||||
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
|
||||
assertEquals(MemberSession.State.BUSY, sessions.get(session.paneId()).orElseThrow().state(),
|
||||
"delivery moves READY → BUSY");
|
||||
|
||||
@@ -173,7 +174,7 @@ class SessionManagerTest {
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
String terminal = session.terminalId();
|
||||
sessions.asPresence().markPresent(terminal);
|
||||
sessions.onDelivered(terminal);
|
||||
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
|
||||
|
||||
sessions.onTurnFailed(terminal);
|
||||
|
||||
@@ -201,7 +202,7 @@ class SessionManagerTest {
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
String terminal = session.terminalId();
|
||||
sessions.asPresence().markPresent(terminal);
|
||||
sessions.onDelivered(terminal);
|
||||
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
|
||||
|
||||
sessions.onTurnFailed(terminal);
|
||||
|
||||
@@ -284,7 +285,7 @@ class SessionManagerTest {
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
String terminal = session.terminalId();
|
||||
sessions.asPresence().markPresent(terminal);
|
||||
sessions.onDelivered(terminal);
|
||||
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
|
||||
|
||||
clock[0] = 100;
|
||||
assertEquals(0, sessions.reapIdle(10), "BUSY session past TTL is never reaped");
|
||||
@@ -301,7 +302,7 @@ class SessionManagerTest {
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
String terminal = session.terminalId();
|
||||
sessions.asPresence().markPresent(terminal);
|
||||
sessions.onDelivered(terminal);
|
||||
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
|
||||
sessions.onTurnComplete(terminal);
|
||||
|
||||
clock[0] = 21;
|
||||
@@ -319,7 +320,7 @@ class SessionManagerTest {
|
||||
MemberSession busy = sessions.acquire("ltms-local", "/busy", "/caller", "owner2");
|
||||
sessions.asPresence().markPresent(ready.terminalId());
|
||||
sessions.asPresence().markPresent(busy.terminalId());
|
||||
sessions.onDelivered(busy.terminalId());
|
||||
sessions.onDelivered(busy.terminalId(), TestTurnTokens.inert(busy.terminalId()));
|
||||
|
||||
clock[0] = 50;
|
||||
assertEquals(1, sessions.reapIdle(30), "only READY past TTL is reaped");
|
||||
@@ -337,9 +338,9 @@ class SessionManagerTest {
|
||||
String terminal = session.terminalId();
|
||||
sessions.asPresence().markPresent(terminal);
|
||||
|
||||
sessions.onDelivered(terminal);
|
||||
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
|
||||
sessions.onTurnComplete(terminal);
|
||||
sessions.onDelivered(terminal);
|
||||
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
|
||||
sessions.onTurnComplete(terminal);
|
||||
|
||||
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
|
||||
@@ -357,13 +358,13 @@ class SessionManagerTest {
|
||||
String terminal = session.terminalId();
|
||||
sessions.asPresence().markPresent(terminal);
|
||||
|
||||
sessions.onDelivered(terminal);
|
||||
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
|
||||
sessions.onTurnComplete(terminal);
|
||||
assertEquals(MemberSession.State.DONE,
|
||||
sessions.get(session.paneId()).orElseThrow().state(),
|
||||
"first turn completes without release");
|
||||
|
||||
sessions.onDelivered(terminal);
|
||||
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
|
||||
sessions.onTurnComplete(terminal);
|
||||
|
||||
assertTrue(sessions.get(session.paneId()).isEmpty(), "session released after cap reached");
|
||||
@@ -379,7 +380,7 @@ class SessionManagerTest {
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
sessions.asPresence().markPresent(session.terminalId());
|
||||
|
||||
sessions.onDelivered(session.terminalId());
|
||||
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
|
||||
assertTrue(sessions.onTurnCompleteWithPostAction(session.terminalId()));
|
||||
|
||||
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
|
||||
@@ -393,7 +394,7 @@ class SessionManagerTest {
|
||||
SessionManager sessions = sessionManager(herdr, () -> 0L, 1, true);
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
sessions.asPresence().markPresent(session.terminalId());
|
||||
sessions.onDelivered(session.terminalId());
|
||||
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
|
||||
|
||||
assertFalse(sessions.hasPostTurnAction(session.terminalId()),
|
||||
"a session at its cap will be released, not reset for reuse");
|
||||
@@ -408,7 +409,7 @@ class SessionManagerTest {
|
||||
SessionManager sessions = sessionManager(herdr, () -> 0L, 0, false);
|
||||
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
|
||||
sessions.asPresence().markPresent(session.terminalId());
|
||||
sessions.onDelivered(session.terminalId());
|
||||
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
|
||||
|
||||
sessions.onTurnComplete(session.terminalId());
|
||||
|
||||
@@ -426,7 +427,7 @@ class SessionManagerTest {
|
||||
MemberSession busy = sessions.acquire("ltms-local", "/busy", "/caller", "ownerB");
|
||||
sessions.asPresence().markPresent(ready.terminalId());
|
||||
sessions.asPresence().markPresent(busy.terminalId());
|
||||
sessions.onDelivered(busy.terminalId());
|
||||
sessions.onDelivered(busy.terminalId(), TestTurnTokens.inert(busy.terminalId()));
|
||||
|
||||
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(100));
|
||||
|
||||
|
||||
@@ -6,9 +6,15 @@ import dev.ltms.bridged.guard.SubscriptionGuard;
|
||||
import dev.ltms.bridged.herdr.AgentControl;
|
||||
import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.LoggerContext;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.bridged.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.bridged.msg.TestTurnTokens;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -178,6 +184,49 @@ class WorktreeSessionManagerTest {
|
||||
assertTrue(sessions.get(paneId).isEmpty(), "released session is no longer retrievable");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-576. A normal {@code COMPLETED} release whose worktree holds uncommitted work must NOT
|
||||
* remove it — {@code --force} would destroy the worker's only copy. The bridge cannot see
|
||||
* uncommitted files, so the worktree is preserved and the release logged at WARN naming the
|
||||
* path, the session, and the cause an operator needs to find the work.
|
||||
*/
|
||||
@Test
|
||||
void releasePreservesDirtyWorktreeAndLogsWarn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
|
||||
.withDirty(true);
|
||||
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
|
||||
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-576", null));
|
||||
|
||||
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
|
||||
ch.qos.logback.classic.Logger sessionLog =
|
||||
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.setContext(ctx);
|
||||
appender.start();
|
||||
sessionLog.addAppender(appender);
|
||||
sessionLog.setLevel(Level.WARN);
|
||||
try {
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertTrue(herdr.called("pane.close"), "release still tears the worker pane down");
|
||||
assertTrue(worktrees.removeCalls().isEmpty(),
|
||||
"a dirty worktree is never removed — it holds the only copy of the work");
|
||||
String warn = appender.list.stream()
|
||||
.filter(e -> e.getLevel().equals(Level.WARN))
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.filter(m -> m.contains("dirty worktree"))
|
||||
.findFirst()
|
||||
.orElse("no dirty-release WARN logged");
|
||||
assertTrue(warn.contains(s.worktree()), "the WARN names the worktree path: " + warn);
|
||||
assertTrue(warn.contains(s.terminalId()), "the WARN names the session: " + warn);
|
||||
assertTrue(warn.contains("COMPLETED"), "the WARN names the release cause: " + warn);
|
||||
} finally {
|
||||
sessionLog.detachAppender(appender);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void drainAllPreservesWorktreeOfIdleSession() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
@@ -204,7 +253,7 @@ class WorktreeSessionManagerTest {
|
||||
new WorktreeRequest("cb-544", null));
|
||||
String terminal = s.terminalId();
|
||||
sessions.asPresence().markPresent(terminal);
|
||||
sessions.onDelivered(terminal); // BUSY, never completes → still BUSY when the timeout hits
|
||||
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal)); // BUSY, never completes → still BUSY when the timeout hits
|
||||
|
||||
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(100));
|
||||
|
||||
|
||||
Reference in New Issue
Block a user