Compare commits
8 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 3b2f395d3d | |||
| 0edc6615fc | |||
| c884802b13 | |||
| 5f5573a24e | |||
| 74b0087ebb | |||
| 927e0151d4 | |||
| 5275922d1d | |||
| e186c7945a |
@@ -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,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,6 +385,9 @@ 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);
|
||||
}
|
||||
// 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.
|
||||
@@ -408,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 {
|
||||
@@ -433,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);
|
||||
@@ -536,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.
|
||||
@@ -551,7 +553,6 @@ public final class MessageService {
|
||||
}
|
||||
} catch (Throwable t) {
|
||||
task.future.completeExceptionally(t);
|
||||
untrackAsyncTarget(task);
|
||||
}
|
||||
});
|
||||
pruneTerminalTickets();
|
||||
@@ -613,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. */
|
||||
@@ -641,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);
|
||||
}
|
||||
@@ -660,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();
|
||||
|
||||
@@ -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
@@ -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