Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 3b2f395d3d | |||
| 0edc6615fc | |||
| c884802b13 | |||
| 5f5573a24e |
@@ -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();
|
||||
@@ -285,12 +287,14 @@ public final class Bridged {
|
||||
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);
|
||||
|
||||
@@ -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<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<>();
|
||||
@@ -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<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();
|
||||
|
||||
@@ -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<String, Set<Object>> byTarget = (Map<String, Set<Object>>) byTargetField.get(messages);
|
||||
Set<Object> 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<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
assertEquals(MessageService.Phase.ASKING, awaitTicketPhase(first, MessageService.Phase.ASKING).phase());
|
||||
|
||||
+22
-2
@@ -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.
|
||||
|
||||
Reference in New Issue
Block a user