From b745e159de7f02c86ab2e6c1506fc36100b65bcd Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 06:27:39 +0200 Subject: [PATCH 1/2] CB-573: carry accepted turn tokens on delivery --- .../main/java/dev/ltms/bridged/Bridged.java | 6 +++--- .../bridged/inject/CompletionResolver.java | 9 +++++---- .../dev/ltms/bridged/inject/Injector.java | 9 +++++---- .../dev/ltms/bridged/inject/TurnListener.java | 4 +++- .../dev/ltms/bridged/msg/MessageService.java | 3 ++- .../java/dev/ltms/bridged/msg/TurnToken.java | 20 +++++++++++++++++++ .../ltms/bridged/session/SessionManager.java | 3 ++- 7 files changed, 40 insertions(+), 14 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/bridged/msg/TurnToken.java diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index fbae9b7..70a10d9 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -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 diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java index e504b96..5259ffd 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java @@ -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 waiter = rendezvous.currentWaiter(target); + void captureBaseline(String target, TurnToken token) { + CompletableFuture waiter = token.waiter(); if (waiter == null) { inFlight.remove(target); // no send is waiting on this delivery — nothing to resolve later return; diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java index 6144c6b..eeb62f3 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java @@ -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 delivered) { + private record Pending(String text, TurnToken token, CompletableFuture delivered) { } /** Per-worker delivery state, guarded by its own monitor (single writer per worker). */ @@ -159,9 +160,9 @@ public final class Injector { *

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 enqueue(String target, String text) { + public CompletableFuture enqueue(String target, String text, TurnToken token) { CompletableFuture 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); } } diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java b/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java index 68e377d..7d453ee 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java @@ -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. */ diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java index 024e274..4c0a151 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -391,13 +391,14 @@ public final class MessageService { 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 delivered = injector.enqueue(target, content); + CompletableFuture 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())); diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/TurnToken.java b/bridged/src/main/java/dev/ltms/bridged/msg/TurnToken.java new file mode 100644 index 0000000..f3f8c2e --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/msg/TurnToken.java @@ -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 waiter; + + public TurnToken(String target, CompletableFuture waiter) { + this.target = target; + this.waiter = waiter; + } + + public String target() { return target; } + public CompletableFuture waiter() { return waiter; } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java index fdb5766..9ebc167 100644 --- a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java +++ b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java @@ -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; @@ -425,7 +426,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) { From 8e2e4c5e73e7e510349975a464d641c9b8e116ba Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 06:36:35 +0200 Subject: [PATCH 2/2] 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. --- .../inject/CompletionResolverTest.java | 8 ++- .../dev/ltms/bridged/inject/InjectorTest.java | 69 ++++++++++--------- .../ltms/bridged/msg/MessageServiceTest.java | 2 +- .../dev/ltms/bridged/msg/TestTurnTokens.java | 20 ++++++ .../bridged/session/SessionManagerTest.java | 31 +++++---- .../session/WorktreeSessionManagerTest.java | 3 +- 6 files changed, 79 insertions(+), 54 deletions(-) create mode 100644 bridged/src/test/java/dev/ltms/bridged/msg/TestTurnTokens.java diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java index 8ae282d..6cf1a40 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/CompletionResolverTest.java @@ -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"); diff --git a/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java b/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java index 61a10b9..fddd39a 100644 --- a/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/inject/InjectorTest.java @@ -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 f = injector.enqueue(T, "hello"); + CompletableFuture 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 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 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 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 f = inj.enqueue(T, "boom"); + CompletableFuture 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 f = injector.enqueue(T, "orphan"); + CompletableFuture 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 delivered = inj.enqueue(T, "delivered"); - CompletableFuture queued = inj.enqueue(T, "queued"); + CompletableFuture delivered = inj.enqueue(T, "delivered", TestTurnTokens.inert(T)); + CompletableFuture 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 forgotten = new ArrayList<>(); Injector inj = new Injector(new AgentControl(herdr), cap, _ -> false, forgotten::add); - CompletableFuture f = inj.enqueue(T, "task"); + CompletableFuture 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 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 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 delivered = inj.enqueue(T, "via-poller"); + CompletableFuture 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 f = inj.enqueue(T, "boom"); + CompletableFuture 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()); diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index a72804e..58e3d04 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -130,7 +130,7 @@ class MessageServiceTest { injector.onStatus(T, AgentStatus.IDLE); // first delivery injector.onStatus(T, AgentStatus.WORKING); // first turn in flight - CompletableFuture queued = injector.enqueue(T, "second task"); + CompletableFuture queued = injector.enqueue(T, "second task", TestTurnTokens.inert(T)); CompletableFuture waiter = rendezvous.currentWaiter(T); injector.drop(T, new HerdrException("agent target sol not found", "agent_not_found", null)); diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/TestTurnTokens.java b/bridged/src/test/java/dev/ltms/bridged/msg/TestTurnTokens.java new file mode 100644 index 0000000..35bee6b --- /dev/null +++ b/bridged/src/test/java/dev/ltms/bridged/msg/TestTurnTokens.java @@ -0,0 +1,20 @@ +package dev.ltms.bridged.msg; + +/** + * Explicit unbound tokens for tests that exercise delivery without an accepted send. + * + *

The waiter is {@code null} on purpose. "No accepted send" is an absence, 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); + } +} diff --git a/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java b/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java index febf37a..422831a 100644 --- a/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/session/SessionManagerTest.java @@ -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)); diff --git a/bridged/src/test/java/dev/ltms/bridged/session/WorktreeSessionManagerTest.java b/bridged/src/test/java/dev/ltms/bridged/session/WorktreeSessionManagerTest.java index 45eb6f8..e0b07e7 100644 --- a/bridged/src/test/java/dev/ltms/bridged/session/WorktreeSessionManagerTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/session/WorktreeSessionManagerTest.java @@ -7,6 +7,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.MemberRole; import org.junit.jupiter.api.Test; @@ -204,7 +205,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));