Compare commits

..

8 Commits

Author SHA1 Message Date
Dai Ha 3b2f395d3d M4: fail tickets on terminal health states
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 57s
2026-08-15 06:32:02 +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 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
10 changed files with 114 additions and 74 deletions
@@ -55,6 +55,7 @@ import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.Executors;
import java.util.concurrent.atomic.AtomicReference;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
@@ -254,6 +255,7 @@ public final class Bridged {
// CB-106: a confirmed turn completion resolves a blocked send whose worker never replied.
Rendezvous rendezvous = new Rendezvous();
CompletionResolver completion = new CompletionResolver(agents, rendezvous);
AtomicReference<MessageService> messagesRef = new AtomicReference<>();
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
MemberPresence presence = sessions.asPresence();
@@ -276,21 +278,23 @@ public final class Bridged {
}
@Override
public void onDelivered(String target, dev.ltms.bridged.msg.TurnToken token) {
completion.onDelivered(target, token);
sessions.onDelivered(target, token);
public void onDelivered(String target) {
completion.onDelivered(target);
sessions.onDelivered(target);
}
@Override
public void onTurnFailed(String target) {
completion.onTurnFailed(target);
sessions.onTurnFailed(target);
failTarget(messagesRef, target, "worker turn failed");
}
@Override
public void onTurnFailed(String target, String reason) {
completion.onTurnFailed(target, reason);
sessions.onTurnFailed(target);
failTarget(messagesRef, target, reason);
}
};
Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads),
@@ -354,6 +358,7 @@ public final class Bridged {
}
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
pushLoop, metrics);
messagesRef.set(messages);
// Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller.
final FleetHealthMonitor healthMonitor;
@@ -361,7 +366,7 @@ public final class Bridged {
Thread.ofVirtual().name("bridge-health-").unstarted(r));
if (cfg.health() != null && cfg.health().isEnabled()) {
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
System::nanoTime, cfg.health().intervalOrDefault());
System::nanoTime, cfg.health().intervalOrDefault(), messages::abandon);
String coverage = FleetHealthMonitor.coverage(true,
cfg.health().notifications() != null && cfg.health().notifications().configured());
if ("detection-only".equals(coverage)) {
@@ -485,6 +490,14 @@ public final class Bridged {
return target -> presence.isPresent(target) || leads.get().containsKey(target);
}
/** Fail outstanding tickets without changing the member lifecycle or worktree state. */
private static void failTarget(AtomicReference<MessageService> messagesRef, String target, String reason) {
MessageService messages = messagesRef.get();
if (messages != null) {
messages.abandon(target, reason);
}
}
/**
* Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504).
*
@@ -16,6 +16,7 @@ import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.function.LongSupplier;
import java.util.function.Supplier;
import java.util.function.BiConsumer;
/** Slow whole-fleet evidence collection. It is deliberately separate from the delivery poller. */
public final class FleetHealthMonitor {
@@ -26,6 +27,7 @@ public final class FleetHealthMonitor {
private final ScheduledExecutorService scheduler;
private final LongSupplier clock;
private final long intervalSeconds;
private final BiConsumer<String, String> failTarget;
private final Map<String, HealthPrior> priors = new HashMap<>();
private final Map<String, HealthState> states = new HashMap<>();
@@ -33,13 +35,20 @@ public final class FleetHealthMonitor {
private static final boolean NOT_YET_OBSERVED = false;
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds) {
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds) {
this(agents, roster, messages, scheduler, clock, intervalSeconds, (_, _) -> { });
}
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds,
BiConsumer<String, String> failTarget) {
this.agents = agents;
this.roster = roster;
this.messages = messages;
this.scheduler = scheduler;
this.clock = clock;
this.intervalSeconds = intervalSeconds;
this.failTarget = failTarget;
}
/** Pure per-member decision seam. */
@@ -69,6 +78,9 @@ public final class FleetHealthMonitor {
HealthDecision decision = decide(snapshot, priors.getOrDefault(session.terminalId(), HealthPrior.NONE),
clock.getAsLong());
priors.put(session.terminalId(), decision.prior());
if (decision.state() == HealthState.GONE || decision.state() == HealthState.NEVER_READY) {
failTarget.accept(session.terminalId(), "fleet health detected " + decision.state());
}
reportTransition(session.terminalId(), decision.state());
}
priors.keySet().retainAll(current);
@@ -2,7 +2,6 @@ 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;
@@ -84,17 +83,17 @@ public final class CompletionResolver implements TurnListener {
}
@Override
public void onDelivered(String target, TurnToken token) {
public void onDelivered(String target) {
// 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, token);
captureBaseline(target);
}
/** Capture the in-flight turn: its waiter and pre-turn baseline (the testable core of {@link #onDelivered}). */
void captureBaseline(String target, TurnToken token) {
CompletableFuture<Rendezvous.Resolution> waiter = token.waiter();
void captureBaseline(String target) {
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
if (waiter == null) {
inFlight.remove(target); // no send is waiting on this delivery — nothing to resolve later
return;
@@ -2,7 +2,6 @@ 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;
@@ -130,7 +129,7 @@ public final class Injector {
}
/** A pending message and the future that completes when it has been delivered. */
private record Pending(String text, TurnToken token, CompletableFuture<Void> delivered) {
private record Pending(String text, CompletableFuture<Void> delivered) {
}
/** Per-worker delivery state, guarded by its own monitor (single writer per worker). */
@@ -160,9 +159,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, TurnToken token) {
public CompletableFuture<Void> enqueue(String target, String text) {
CompletableFuture<Void> delivered = new CompletableFuture<>();
Pending p = new Pending(text, token, delivered);
Pending p = new Pending(text, 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
@@ -357,7 +356,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, sent.token());
turnListener.onDelivered(target);
sent.delivered().complete(null);
}
}
@@ -1,7 +1,5 @@
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
@@ -59,7 +57,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, TurnToken token) {
default void onDelivered(String target) {
}
/** 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,14 +385,16 @@ public final class MessageService {
// failed send leaves no stale waiter behind.
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target);
try {
TurnToken token = new TurnToken(target, reply);
if (task != null) {
asyncTasksByWaiter.put(reply, task);
}
// 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, token);
CompletableFuture<Void> delivered = injector.enqueue(target, content);
try {
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId()));
@@ -409,6 +411,7 @@ public final class MessageService {
throw new IllegalStateException("interrupted awaiting reply from " + target, e);
}
} finally {
asyncTasksByWaiter.remove(reply);
rendezvous.close(target, reply);
}
} finally {
@@ -434,11 +437,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);
@@ -537,13 +544,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.
@@ -552,7 +553,6 @@ public final class MessageService {
}
} catch (Throwable t) {
task.future.completeExceptionally(t);
untrackAsyncTarget(task);
}
});
pruneTerminalTickets();
@@ -614,17 +614,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. */
@@ -642,7 +639,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);
}
@@ -661,17 +657,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();
@@ -1,20 +0,0 @@
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; }
}
@@ -4,7 +4,6 @@ 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;
@@ -426,7 +425,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, TurnToken token) {
public void onDelivered(String target) {
MemberSession current = findByTerminal(target);
if (current == null) return;
if (current.state() != MemberSession.State.READY && current.state() != MemberSession.State.DONE) {
@@ -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());
+22 -2
View File
@@ -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.