From c884802b1312fae4f87e909440da92b7ac29211d Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 06:28:53 +0200 Subject: [PATCH 1/2] CB-577: remove obsolete async target tracking --- .../dev/ltms/bridged/msg/MessageService.java | 25 +------------------ .../ltms/bridged/msg/MessageServiceTest.java | 15 ----------- 2 files changed, 1 insertion(+), 39 deletions(-) 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..ca75d04 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -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,6 @@ public final class MessageService { private final Metrics metrics; // CB-502: nullable — no registry in unit tests private final ConcurrentHashMap sessionLocks = new ConcurrentHashMap<>(); private final ConcurrentHashMap tasks = new ConcurrentHashMap<>(); - /** Async tasks that have accepted delivery for a target. */ - private final ConcurrentHashMap> asyncTasksByTarget = new ConcurrentHashMap<>(); /** Async task that owns each exact forward rendezvous waiter. */ private final ConcurrentHashMap, Task> asyncTasksByWaiter = new ConcurrentHashMap<>(); @@ -547,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. @@ -562,7 +553,6 @@ public final class MessageService { } } catch (Throwable t) { task.future.completeExceptionally(t); - untrackAsyncTarget(task); } }); pruneTerminalTickets(); @@ -642,7 +632,6 @@ public final class MessageService { if (forgetTurn) { asyncTasksByTurn.remove(turnId, task); task.turnId = null; - untrackAsyncTarget(task); } } } @@ -650,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); } @@ -669,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 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(); 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..7f25043 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -10,9 +10,6 @@ import dev.ltms.bridged.inject.Injector; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; -import java.lang.reflect.Field; -import java.util.Map; -import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; @@ -697,22 +694,10 @@ class MessageServiceTest { } @Test - @SuppressWarnings("unchecked") void asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter() throws Exception { String first = messages.sendAsync(T, "first task"); awaitWaiting(); - // Model the resolveQuestion/markAsyncQuestion race: another accepted task reached the target set. - Field byTargetField = MessageService.class.getDeclaredField("asyncTasksByTarget"); - byTargetField.setAccessible(true); - Map> byTarget = (Map>) byTargetField.get(messages); - Set targetTasks = byTarget.get(T); - Class taskClass = Class.forName(MessageService.class.getName() + "$Task"); - var constructor = taskClass.getDeclaredConstructor(String.class); - constructor.setAccessible(true); - targetTasks.clear(); - targetTasks.add(constructor.newInstance(T)); - CompletableFuture ask = CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000)); assertEquals(MessageService.Phase.ASKING, awaitTicketPhase(first, MessageService.Phase.ASKING).phase()); From 0edc6615fcd226a1d8702964e78a0e7e2716393d Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 06:29:27 +0200 Subject: [PATCH 2/2] =?UTF-8?q?M4:=20correct=20unit=202=20criterion=201=20?= =?UTF-8?q?=E2=80=94=20TurnToken=20cannot=20carry=20the=20session=20turn?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- docs/M4-Fleet-Health.md | 24 ++++++++++++++++++++++-- 1 file changed, 22 insertions(+), 2 deletions(-) diff --git a/docs/M4-Fleet-Health.md b/docs/M4-Fleet-Health.md index 11e36c2..f954a99 100644 --- a/docs/M4-Fleet-Health.md +++ b/docs/M4-Fleet-Health.md @@ -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.