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();
|
||||
@@ -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
@@ -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