Compare commits

...

15 Commits

Author SHA1 Message Date
Dai Ha c393600921 CB-576: prove teardown survives an already-gone worktree 2026-08-15 09:31:10 +02:00
Dai Ha b525b0f08f CB-576: hasUncommitted tolerates an already-gone worktree
CI / contract (pull_request) Successful in 1m8s
CI / build (pull_request) Successful in 1m37s
2026-08-15 08:47:12 +02:00
Dai Ha 9118ce2537 CB-576: release preserves a dirty worktree instead of deleting it
CI / contract (pull_request) Successful in 1m2s
CI / build (pull_request) Successful in 1m35s
2026-08-15 07:34:06 +02:00
Dai Ha 2f48e08f1f Merge CB-577 follow-up: drop the target-keyed async index
CI / contract (push) Successful in 43s
CI / build (push) Successful in 55s
asyncTasksByWaiter correlates an async question by the exact rendezvous
waiter, so the target-keyed set it replaced can no longer decide
anything. Keeping it meant two indexes of the same fact, one of them
ambiguous whenever a target has two accepted tickets.

The race the old test modelled by reflection is gone with it: identity
keys make 'some other task reached this target' unrepresentable, so
there is no longer a wrong task for the question to land on.

Also carries the criterion-1 doc correction, which is identical to
aac29d6 on a different parent.
2026-08-15 06:40:31 +02:00
Dai Ha fec284e7cb Merge M4 unit 2a: accepted-turn identity carried to delivery
CI / build (push) Successful in 59s
CI / contract (push) Successful in 1m17s
TurnToken, owned by MessageService, binds a target to the exact
rendezvous waiter for one accepted send. Injector.Pending carries it and
the delivery callback hands it to CompletionResolver, so the baseline is
bound to the send it belongs to by construction rather than by a lookup
that could pick a different one.

The callback signature is required, not a defaulted overload: a delivery
with no token is exactly the unbound baseline this unit forbids, so a
default would let a caller silently produce it.

The token deliberately omits the session turn number. MessageService
owns acceptance but never learns of delivery, and
CompletionResolver.onDelivered runs before SessionManager.onDelivered,
so the number does not exist yet at the only point the token could
capture it. docs/M4-Fleet-Health.md criterion 1 records this and the two
rejected alternatives.

Still open for the next slice: the missing/post-restart baseline test and
the no-replay test.
2026-08-15 06:36:46 +02:00
Dai Ha 8e2e4c5e73 M4 unit 2a: migrate test call sites to the required turn token
The delivery callback now requires a TurnToken, so 54 test call sites
had to pass one. They use an explicit TestTurnTokens.inert(target)
rather than a defaulted overload, because a delivery with no token is
the unbound baseline this unit forbids.

The first version of inert() returned a fresh CompletableFuture as the
waiter, which turned captureBaselineSkipsTheReadWhenNoSendIsWaiting red:
the resolver saw a non-null waiter, concluded a turn was in flight, and
scraped a pane no send was blocked on. An inert value must omit the
fact, not invent it, so the waiter is now null and the production skip
fires as designed.
2026-08-15 06:36:35 +02:00
Dai Ha 0edc6615fc M4: correct unit 2 criterion 1 — TurnToken cannot carry the session turn
The criterion required the token to bind the session turn number. Three
independent refusals from the implementer showed why that is not
implementable at this layer: MessageService owns acceptance but never
learns of delivery, and CompletionResolver.onDelivered runs before
SessionManager.onDelivered, so the turn number does not exist yet at the
only point the token could capture it.

Records both rejected alternatives and why, so the next reader does not
re-derive them: a target-keyed registry restores the ambiguity the token
exists to remove, and injecting a turn counter couples layers to fill a
field nothing reads yet.
2026-08-15 06:32:02 +02:00
Dai Ha b745e159de CB-573: carry accepted turn tokens on delivery 2026-08-15 06:30:41 +02:00
Dai Ha aac29d604c M4: correct unit 2 criterion 1 — TurnToken cannot carry the session turn
CI / contract (push) Successful in 1m0s
CI / build (push) Successful in 1m32s
The criterion required the token to bind the session turn number. Three
independent refusals from the implementer showed why that is not
implementable at this layer: MessageService owns acceptance but never
learns of delivery, and CompletionResolver.onDelivered runs before
SessionManager.onDelivered, so the turn number does not exist yet at the
only point the token could capture it.

Records both rejected alternatives and why, so the next reader does not
re-derive them: a target-keyed registry restores the ambiguity the token
exists to remove, and injecting a turn counter couples layers to fill a
field nothing reads yet.
2026-08-15 06:29:27 +02:00
Dai Ha c884802b13 CB-577: remove obsolete async target tracking
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 56s
2026-08-15 06:28:53 +02:00
Dai Ha 5f5573a24e Merge CB-577: correlate an async question by its exact waiter
CI / contract (push) Failing after 0s
CI / build (push) Successful in 1m13s
markAsyncQuestion picked the first not-done task out of an unordered
set, so between resolveQuestion waking the first async send and the
question being recorded, a queued second send could join the set and
take the question. A lead answering with bridge_send{turnId} would then
resume a turn it did not mean to.

Each async task is now indexed by its exact rendezvous waiter, which has
identity semantics, so no other task can hold the same key. The question
is recorded before resolveQuestion, with a rollback when no waiter is
there, which closes the window rather than narrowing it.

An unanswered async question stays PENDING — the worker resumes after
its ask times out, so the delegation is not failed — and its stale
target tracking is now cleared instead of leaking.
2026-08-15 06:27:16 +02:00
Dai Ha 74b0087ebb CB-577: model async question ownership race
CI / build (pull_request) Successful in 1m29s
CI / contract (pull_request) Failing after 0s
2026-08-15 06:24:05 +02:00
Dai Ha 927e0151d4 CB-577: test async question waiter ownership 2026-08-15 06:23:22 +02:00
Dai Ha 5275922d1d CB-577: handle questions without async waiters 2026-08-15 06:22:22 +02:00
Dai Ha e186c7945a CB-577: correlate async questions to turns 2026-08-15 06:21:53 +02:00
18 changed files with 376 additions and 102 deletions
@@ -276,9 +276,9 @@ public final class Bridged {
}
@Override
public void onDelivered(String target) {
completion.onDelivered(target);
sessions.onDelivered(target);
public void onDelivered(String target, dev.ltms.bridged.msg.TurnToken token) {
completion.onDelivered(target, token);
sessions.onDelivered(target, token);
}
@Override
@@ -2,6 +2,7 @@ package dev.ltms.bridged.inject;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.msg.TurnToken;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -83,17 +84,17 @@ public final class CompletionResolver implements TurnListener {
}
@Override
public void onDelivered(String target) {
public void onDelivered(String target, TurnToken token) {
// Capture the exact waiter this turn belongs to (CB-116) and snapshot the pane's pre-turn
// content — what it shows *before* the just-delivered turn produces output — as the staleness
// reference (CB-115). Done synchronously (like the delivering send itself) so both are in
// place before this turn's completion can fire.
captureBaseline(target);
captureBaseline(target, token);
}
/** Capture the in-flight turn: its waiter and pre-turn baseline (the testable core of {@link #onDelivered}). */
void captureBaseline(String target) {
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
void captureBaseline(String target, TurnToken token) {
CompletableFuture<Rendezvous.Resolution> waiter = token.waiter();
if (waiter == null) {
inFlight.remove(target); // no send is waiting on this delivery — nothing to resolve later
return;
@@ -2,6 +2,7 @@ package dev.ltms.bridged.inject;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.msg.TurnToken;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -129,7 +130,7 @@ public final class Injector {
}
/** A pending message and the future that completes when it has been delivered. */
private record Pending(String text, CompletableFuture<Void> delivered) {
private record Pending(String text, TurnToken token, CompletableFuture<Void> delivered) {
}
/** Per-worker delivery state, guarded by its own monitor (single writer per worker). */
@@ -159,9 +160,9 @@ public final class Injector {
* <p>Uses an atomic map update so a concurrent {@link #drop} cannot slip between "find the
* target" and "queue the message" and orphan it in a target it just removed.
*/
public CompletableFuture<Void> enqueue(String target, String text) {
public CompletableFuture<Void> enqueue(String target, String text, TurnToken token) {
CompletableFuture<Void> delivered = new CompletableFuture<>();
Pending p = new Pending(text, delivered);
Pending p = new Pending(text, token, delivered);
targets.compute(target, (_, existing) -> {
Target t = (existing != null) ? existing : new Target();
t.add(p); // synchronized on the Target monitor — atomic with a concurrent drop
@@ -356,7 +357,7 @@ public final class Injector {
} else {
// Baseline the pane's pre-turn content so a misattributed completion (no new output)
// can't resolve this send with the previous turn's stale answer (CB-115).
turnListener.onDelivered(target);
turnListener.onDelivered(target, sent.token());
sent.delivered().complete(null);
}
}
@@ -1,5 +1,7 @@
package dev.ltms.bridged.inject;
import dev.ltms.bridged.msg.TurnToken;
/**
* Notified when a worker's delegated turn is observed to complete — a confirmed
* {@code WORKING → IDLE} transition after a delivery. This is the CB-106 completion signal the
@@ -57,7 +59,7 @@ public interface TurnListener {
* resolve the send with the previous turn's stale answer. A default no-op keeps the interface
* functional for callers that don't scrape.
*/
default void onDelivered(String target) {
default void onDelivered(String target, TurnToken token) {
}
/** No-op default for callers that only need delivery, not completion signalling. */
@@ -9,7 +9,6 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
@@ -171,8 +170,9 @@ public final class MessageService {
private final Metrics metrics; // CB-502: nullable — no registry in unit tests
private final ConcurrentHashMap<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Task> tasks = new ConcurrentHashMap<>();
/** Async tasks that have accepted delivery for a target. */
private final ConcurrentHashMap<String, Set<Task>> asyncTasksByTarget = new ConcurrentHashMap<>();
/** Async task that owns each exact forward rendezvous waiter. */
private final ConcurrentHashMap<CompletableFuture<Rendezvous.Resolution>, Task> asyncTasksByWaiter =
new ConcurrentHashMap<>();
/** Async tickets paused on a specific {@code bridge_ask} turn. */
private final ConcurrentHashMap<String, Task> asyncTasksByTurn = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
@@ -385,13 +385,17 @@ public final class MessageService {
// failed send leaves no stale waiter behind.
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target);
try {
if (task != null) {
asyncTasksByWaiter.put(reply, task);
}
TurnToken token = new TurnToken(target, reply);
// The send has won the lock; the accepted-delivery hook records delegator ownership
// here (CB-548). It runs BEFORE enqueue so a throwing hook — onAccepted is now a
// public callback — fails the send without queuing a message that would orphan.
if (onAccepted != null) {
onAccepted.run();
}
CompletableFuture<Void> delivered = injector.enqueue(target, content);
CompletableFuture<Void> delivered = injector.enqueue(target, content, token);
try {
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId()));
@@ -408,6 +412,7 @@ public final class MessageService {
throw new IllegalStateException("interrupted awaiting reply from " + target, e);
}
} finally {
asyncTasksByWaiter.remove(reply);
rendezvous.close(target, reply);
}
} finally {
@@ -433,11 +438,15 @@ public final class MessageService {
if (ticket.fresh()) {
// Register the reverse waiter first, then surface the question — so the answer, which can
// arrive the instant the primary reacts, always finds an open waiter to resolve.
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(workerSession);
Task task = markAsyncQuestion(waiter, question, ticket.turnId());
if (!rendezvous.resolveQuestion(workerSession, question, ticket.turnId())) {
if (task != null) {
clearAsyncQuestion(ticket.turnId(), true);
}
rendezvous.closeAsk(ticket.turnId());
return new AskResult(AskOutcome.NO_WAITER, null); // no primary is blocked on this worker
}
markAsyncQuestion(workerSession, question, ticket.turnId());
}
try {
String answer = ticket.answer().get(timeoutMillis, TimeUnit.MILLISECONDS);
@@ -536,13 +545,7 @@ public final class MessageService {
tasks.put(ticket, task);
asyncExecutor.submit(() -> {
try {
Runnable trackingAccepted = () -> {
if (onAccepted != null) {
onAccepted.run();
}
asyncTasksByTarget.computeIfAbsent(target, _ -> ConcurrentHashMap.newKeySet()).add(task);
};
Reply result = send(target, content, ASYNC_TIMEOUT_MS, trackingAccepted, task);
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task);
if (result.outcome() == Outcome.QUESTION) {
// Keep the accepted owner until answer() finishes it. markAsyncQuestion may run
// just after resolveQuestion wakes this thread.
@@ -551,7 +554,6 @@ public final class MessageService {
}
} catch (Throwable t) {
task.future.completeExceptionally(t);
untrackAsyncTarget(task);
}
});
pruneTerminalTickets();
@@ -613,17 +615,14 @@ public final class MessageService {
}
/** Record the active question for an async ticket; blocking sends have no entry and stay unchanged. */
private void markAsyncQuestion(String target, String text, String turnId) {
Set<Task> targetTasks = asyncTasksByTarget.get(target);
Task task = targetTasks == null ? null : targetTasks.stream()
.filter(candidate -> !candidate.future.isDone())
.findFirst()
.orElse(null);
private Task markAsyncQuestion(CompletableFuture<Rendezvous.Resolution> waiter, String text, String turnId) {
Task task = waiter == null ? null : asyncTasksByWaiter.get(waiter);
if (task != null) {
task.question = new Reply(Outcome.QUESTION, text, turnId);
task.turnId = turnId;
asyncTasksByTurn.put(turnId, task);
}
return task;
}
/** Clear an answered or lapsed question, but only when it matches the ticket's current turn. */
@@ -641,7 +640,6 @@ public final class MessageService {
/** Complete and detach an async ticket after its worker's actual terminal reply. */
private void finishAsyncTask(Task task, Reply result) {
task.future.complete(result);
untrackAsyncTarget(task);
if (task.turnId != null) {
asyncTasksByTurn.remove(task.turnId, task);
}
@@ -660,17 +658,6 @@ public final class MessageService {
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
}
/** Stop tracking a task once it no longer owns an accepted target turn. */
private void untrackAsyncTarget(Task task) {
Set<Task> targetTasks = asyncTasksByTarget.get(task.target);
if (targetTasks != null) {
targetTasks.remove(task);
if (targetTasks.isEmpty()) {
asyncTasksByTarget.remove(task.target, targetTasks);
}
}
}
/** Release the async executor. */
public void close() {
asyncExecutor.shutdown();
@@ -0,0 +1,20 @@
package dev.ltms.bridged.msg;
import java.util.concurrent.CompletableFuture;
/**
* Identity for one accepted send. The session turn is deliberately absent: CompletionResolver's
* delivery callback runs before SessionManager.onDelivered, so binding it needs a later ordering design.
*/
public final class TurnToken {
private final String target;
private final CompletableFuture<Rendezvous.Resolution> waiter;
public TurnToken(String target, CompletableFuture<Rendezvous.Resolution> waiter) {
this.target = target;
this.waiter = waiter;
}
public String target() { return target; }
public CompletableFuture<Rendezvous.Resolution> waiter() { return waiter; }
}
@@ -167,6 +167,22 @@ public final class GitWorktrees implements Worktrees {
exec("git", "-C", repoRoot, "worktree", "remove", "--force", worktreePath);
}
@Override
public boolean hasUncommitted(String worktreePath) {
// A worktree that is already gone holds no work to lose, and it must not break teardown:
// git -C <missing-dir> status exits non-zero and would throw where release() is mid-way
// through stopping a pane. Mirror remove()'s already-gone tolerance by treating it as clean.
Path p = Path.of(worktreePath);
if (!Files.exists(p)) {
log.debug("worktree {} already gone — nothing can be uncommitted", worktreePath);
return false;
}
// No --untracked-files=no: the exact shape of the work lost in CB-576 was a new file
// that was never added, so an untracked-only worktree is still dirty.
String out = exec("git", "-C", worktreePath, "status", "--porcelain");
return !out.isBlank();
}
@Override
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
if (overlay == null || overlay.isEmpty()) {
@@ -4,6 +4,7 @@ import dev.ltms.bridged.auth.MemberLifecycle;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.inject.TurnListener;
import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.msg.TurnToken;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.peer.PeerHandle;
import dev.ltms.bridged.peer.PeerLauncher;
@@ -200,6 +201,16 @@ public final class SessionManager implements TurnListener {
removed.paneId(), removed.terminalId(), removed.state(), cause);
if (preserveWorktree && removed.worktree() != null) {
logPreservedForShutdown(removed);
} else if (removed.worktree() != null && worktrees.hasUncommitted(removed.worktree())) {
// CB-576: a release that would otherwise remove the worktree finds it holding
// uncommitted work the bridge cannot see. A worker that ends a turn without
// committing (normally because it stopped to ask a question or refused the turn)
// has its only copy of that work in the worktree. Remove would --force-delete it,
// so preserve the directory and tell an operator where to find it.
preserveWorktree = true;
log.warn("release {} preserves dirty worktree {} for pane={} terminal={}: "
+ "the worktree holds uncommitted changes that --force remove would destroy",
cause, removed.worktree(), removed.paneId(), removed.terminalId());
}
// CB-516: a send still waiting on this worker can never be answered now. Tell the
// listener BEFORE the pane is torn down, so a blocked caller fails fast with a real
@@ -425,7 +436,7 @@ public final class SessionManager implements TurnListener {
* can be re-delivered for multi-turn reuse until it is released.
*/
@Override
public void onDelivered(String target) {
public void onDelivered(String target, TurnToken token) {
MemberSession current = findByTerminal(target);
if (current == null) return;
if (current.state() != MemberSession.State.READY && current.state() != MemberSession.State.DONE) {
@@ -10,6 +10,18 @@ public interface Worktrees {
/** git -C <repoRoot> worktree remove --force <path>. Idempotent (already-gone tolerated). */
void remove(String repoRoot, String worktreePath);
/**
* True when the worktree holds uncommitted changes the bridge cannot see: tracked
* modifications, staged files, or untracked files. {@code git status --porcelain} is the
* test; an empty result means clean. Callers use this to decide whether removing the
* worktree would silently destroy a worker's only copy of its work.
*
* <p>An already-gone worktree is reported as clean (no throw), matching {@link #remove}'s
* idempotent contract: a path that does not exist holds no work to lose, and must not break
* a teardown that is mid-way through stopping the pane.
*/
boolean hasUncommitted(String worktreePath);
/** Copy each existing overlay path repoRoot→worktree; mark tracked ones --skip-worktree. */
void overlayParity(String repoRoot, String worktreePath, List<String> overlay);
@@ -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
@@ -285,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
@@ -300,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
@@ -315,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());
@@ -324,7 +325,7 @@ 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");
}
@@ -333,8 +334,8 @@ class InjectorTest {
void dropPassesTheRealCauseForQueuedAndDeliveredWork() {
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), cap);
CompletableFuture<Void> delivered = inj.enqueue(T, "delivered");
CompletableFuture<Void> queued = inj.enqueue(T, "queued");
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
@@ -352,7 +353,7 @@ class InjectorTest {
// 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
@@ -367,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));
@@ -386,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);
@@ -414,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);
@@ -439,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());
@@ -455,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");
}
@@ -476,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()
@@ -500,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();
@@ -519,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());
@@ -127,7 +127,7 @@ class MessageServiceTest {
injector.onStatus(T, AgentStatus.IDLE); // first delivery
injector.onStatus(T, AgentStatus.WORKING); // first turn in flight
CompletableFuture<Void> queued = injector.enqueue(T, "second task");
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));
@@ -676,6 +676,41 @@ class MessageServiceTest {
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());
@@ -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,8 +29,11 @@ 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";
/** Worktree paths that currently exist, mirroring real {@code Files.exists} for the gone case. */
private final Set<String> worktreePaths = ConcurrentHashMap.newKeySet();
public FakeWorktrees withRepoRoot(String root) {
this.repoRoot = root;
@@ -61,6 +64,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));
@@ -69,7 +78,15 @@ public final class FakeWorktrees implements Worktrees {
}
// The branch already carries a unique nonce, so the derived path is distinct per acquire
// without an extra counter — keep it a pure function of the branch the test can predict.
return prefix + "/" + branch.replace('/', '_');
String path = prefix + "/" + branch.replace('/', '_');
worktreePaths.add(path);
return path;
}
/** Model an operator / {@code git worktree prune} removing the worktree before release. */
public FakeWorktrees markGone(String worktreePath) {
worktreePaths.remove(worktreePath);
return this;
}
@Override
@@ -77,6 +94,16 @@ public final class FakeWorktrees implements Worktrees {
removeCalls.add(new RemoveCall(repoRoot, worktreePath));
}
@Override
public boolean hasUncommitted(String worktreePath) {
// A path that does not exist (never added, or marked gone) is reported clean, mirroring
// GitWorktrees' already-gone guard — never an error, so teardown still completes.
if (!worktreePaths.contains(worktreePath)) {
return false;
}
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.PeerUnreachableException;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
@@ -101,7 +102,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));
@@ -121,7 +122,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");
@@ -153,7 +154,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);
@@ -181,7 +182,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);
@@ -264,7 +265,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");
@@ -281,7 +282,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;
@@ -299,7 +300,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");
@@ -317,9 +318,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();
@@ -337,13 +338,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");
@@ -359,7 +360,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();
@@ -373,7 +374,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");
@@ -388,7 +389,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());
@@ -406,7 +407,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,14 +6,21 @@ 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;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.*;
@@ -178,6 +185,77 @@ 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);
}
}
/**
* CB-576 review. A worktree that is already gone (operator cleanup, {@code git worktree prune},
* an earlier half-completed release) must not break teardown. {@code hasUncommitted} reports the
* missing path clean, so release still runs {@code notifyReleased} (the CB-516 fast-fail for a
* blocked {@code bridge_send} caller) and {@code launcher.stop} (so the pane is not orphaned),
* and falls through to the already-gone-tolerant {@code remove}.
*/
@Test
void releaseStillStopsPaneAndNotifiesWhenWorktreeIsGone() {
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-576g", null));
AtomicReference<String> releasedTerminal = new AtomicReference<>();
sessions.onRelease(releasedTerminal::set);
worktrees.markGone(s.worktree());
sessions.release(s.paneId());
assertEquals(s.terminalId(), releasedTerminal.get(),
"notifyReleased must still fire when the worktree is already gone (CB-516)");
assertTrue(herdr.called("pane.close"),
"the pane must still be stopped when the worktree is already gone");
assertEquals(1, worktrees.removeCalls().size(),
"release still calls the already-gone-tolerant remove");
}
@Test
void drainAllPreservesWorktreeOfIdleSession() {
FakeHerdr herdr = new FakeHerdr();
@@ -204,7 +282,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));
+23 -3
View File
@@ -237,7 +237,7 @@ state never presents stop as the only action.
| Release cause | Process action | Provisioned worktree |
|---|---|---|
| `SPAWN_ROLLBACK` before registration or delivery | Stop and clean up | Remove |
| `COMPLETED` for `READY` or `DONE` without pending work, idle TTL, or successful context-cap completion | Stop | Remove under completed policy |
| `COMPLETED` for `READY` or `DONE` without pending work, idle TTL, or successful context-cap completion | Stop | Remove only if clean; preserve a dirty worktree (CB-576) |
| `NEVER_READY` | Stop | Preserve |
| `GONE` | Best-effort stop | Preserve |
| `TURN_FAILED` or lead abort while `BUSY` or `FAILED` | Stop | Preserve |
@@ -744,8 +744,28 @@ worktree discovery.
Acceptance criteria:
1. Every accepted send receives a stable `TurnToken` tied to target, exact waiter, session turn,
delivery baseline, and task outcome.
1. Every accepted send receives a stable `TurnToken` tied to target, exact waiter, and delivery
baseline.
**Corrected during implementation (2026-08-15).** This criterion first also required the session
turn number and the task outcome. That is not implementable at this layer, and the implementer
refused it three times rather than fabricate a value — correctly. The reason is an ordering fact
that is invisible from any single class: `MessageService` owns acceptance and holds the waiter and
the async `Task`, but it learns nothing about delivery, because the delivery event goes to
`CompletionResolver` through `TurnListener.onDelivered`. And `CompletionResolver.onDelivered` runs
*before* `SessionManager.onDelivered`, so the session turn number does not exist yet at the only
point where the token could capture it.
Two ways out were rejected. A shared registry keyed by target reintroduces exactly the "whichever
send happens to be waiting" ambiguity the token exists to remove — the same weak claim
`Rendezvous.currentWaiter` warns about. Injecting a turn counter into `MessageService` adds a
required cross-layer dependency to populate a field that nothing in this slice reads, which is
speculative coupling across a boundary already shown to be fragile.
So the token identifies the **accepted send**, and `SessionManager` keeps verifying its own
delivery separately. Repair (criterion 2) does need the session turn; binding it means resolving
that acceptance-versus-delivery ordering first, and that work belongs to the repair unit, not
here. The token record carries a comment saying the field is deliberately absent.
2. Repair requires the same `BUSY` token, two raw `IDLE` or `DONE` snapshots, no conflicting
observation, exact open waiter, successful baseline, and new recognised assistant output.
3. Missing, failed, late, or post-restart baseline never authorises repair.