Compare commits

...

1 Commits

Author SHA1 Message Date
Dai Ha 38f13f4b3c CB-573: carry accepted turn tokens on delivery
CI / contract (pull_request) Failing after 40s
CI / build (pull_request) Failing after 1m9s
2026-08-15 06:27:39 +02:00
7 changed files with 40 additions and 14 deletions
@@ -276,9 +276,9 @@ public final class Bridged {
} }
@Override @Override
public void onDelivered(String target) { public void onDelivered(String target, dev.ltms.bridged.msg.TurnToken token) {
completion.onDelivered(target); completion.onDelivered(target, token);
sessions.onDelivered(target); sessions.onDelivered(target, token);
} }
@Override @Override
@@ -2,6 +2,7 @@ package dev.ltms.bridged.inject;
import dev.ltms.bridged.herdr.AgentControl; import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.msg.Rendezvous; import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.msg.TurnToken;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
@@ -83,17 +84,17 @@ public final class CompletionResolver implements TurnListener {
} }
@Override @Override
public void onDelivered(String target) { public void onDelivered(String target, TurnToken token) {
// Capture the exact waiter this turn belongs to (CB-116) and snapshot the pane's pre-turn // 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 // 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 // reference (CB-115). Done synchronously (like the delivering send itself) so both are in
// place before this turn's completion can fire. // place before this turn's completion can fire.
captureBaseline(target); captureBaseline(target, token);
} }
/** Capture the in-flight turn: its waiter and pre-turn baseline (the testable core of {@link #onDelivered}). */ /** Capture the in-flight turn: its waiter and pre-turn baseline (the testable core of {@link #onDelivered}). */
void captureBaseline(String target) { void captureBaseline(String target, TurnToken token) {
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target); CompletableFuture<Rendezvous.Resolution> waiter = token.waiter();
if (waiter == null) { if (waiter == null) {
inFlight.remove(target); // no send is waiting on this delivery — nothing to resolve later inFlight.remove(target); // no send is waiting on this delivery — nothing to resolve later
return; return;
@@ -2,6 +2,7 @@ package dev.ltms.bridged.inject;
import dev.ltms.bridged.herdr.AgentControl; import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus; import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.msg.TurnToken;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
@@ -129,7 +130,7 @@ public final class Injector {
} }
/** A pending message and the future that completes when it has been delivered. */ /** A pending message and the future that completes when it has been delivered. */
private record Pending(String text, CompletableFuture<Void> delivered) { private record Pending(String text, TurnToken token, CompletableFuture<Void> delivered) {
} }
/** Per-worker delivery state, guarded by its own monitor (single writer per worker). */ /** Per-worker delivery state, guarded by its own monitor (single writer per worker). */
@@ -159,9 +160,9 @@ public final class Injector {
* <p>Uses an atomic map update so a concurrent {@link #drop} cannot slip between "find the * <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. * target" and "queue the message" and orphan it in a target it just removed.
*/ */
public CompletableFuture<Void> enqueue(String target, String text) { public CompletableFuture<Void> enqueue(String target, String text, TurnToken token) {
CompletableFuture<Void> delivered = new CompletableFuture<>(); CompletableFuture<Void> delivered = new CompletableFuture<>();
Pending p = new Pending(text, delivered); Pending p = new Pending(text, token, delivered);
targets.compute(target, (_, existing) -> { targets.compute(target, (_, existing) -> {
Target t = (existing != null) ? existing : new Target(); Target t = (existing != null) ? existing : new Target();
t.add(p); // synchronized on the Target monitor — atomic with a concurrent drop t.add(p); // synchronized on the Target monitor — atomic with a concurrent drop
@@ -356,7 +357,7 @@ public final class Injector {
} else { } else {
// Baseline the pane's pre-turn content so a misattributed completion (no new output) // 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). // can't resolve this send with the previous turn's stale answer (CB-115).
turnListener.onDelivered(target); turnListener.onDelivered(target, sent.token());
sent.delivered().complete(null); sent.delivered().complete(null);
} }
} }
@@ -1,5 +1,7 @@
package dev.ltms.bridged.inject; package dev.ltms.bridged.inject;
import dev.ltms.bridged.msg.TurnToken;
/** /**
* Notified when a worker's delegated turn is observed to complete — a confirmed * 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 * {@code WORKING → IDLE} transition after a delivery. This is the CB-106 completion signal the
@@ -57,7 +59,7 @@ public interface TurnListener {
* resolve the send with the previous turn's stale answer. A default no-op keeps the interface * resolve the send with the previous turn's stale answer. A default no-op keeps the interface
* functional for callers that don't scrape. * functional for callers that don't scrape.
*/ */
default void onDelivered(String target) { default void onDelivered(String target, TurnToken token) {
} }
/** No-op default for callers that only need delivery, not completion signalling. */ /** No-op default for callers that only need delivery, not completion signalling. */
@@ -385,13 +385,14 @@ public final class MessageService {
// failed send leaves no stale waiter behind. // failed send leaves no stale waiter behind.
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target); CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target);
try { try {
TurnToken token = new TurnToken(target, reply);
// The send has won the lock; the accepted-delivery hook records delegator ownership // 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 // 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. // public callback — fails the send without queuing a message that would orphan.
if (onAccepted != null) { if (onAccepted != null) {
onAccepted.run(); onAccepted.run();
} }
CompletableFuture<Void> delivered = injector.enqueue(target, content); CompletableFuture<Void> delivered = injector.enqueue(target, content, token);
try { try {
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS); Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId())); return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId()));
@@ -0,0 +1,20 @@
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,6 +4,7 @@ import dev.ltms.bridged.auth.MemberLifecycle;
import dev.ltms.bridged.herdr.Agent; import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.inject.TurnListener; import dev.ltms.bridged.inject.TurnListener;
import dev.ltms.bridged.inject.MemberPresence; import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.msg.TurnToken;
import dev.ltms.bridged.peer.MemberRole; import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.peer.PeerHandle; import dev.ltms.bridged.peer.PeerHandle;
import dev.ltms.bridged.peer.PeerLauncher; import dev.ltms.bridged.peer.PeerLauncher;
@@ -425,7 +426,7 @@ public final class SessionManager implements TurnListener {
* can be re-delivered for multi-turn reuse until it is released. * can be re-delivered for multi-turn reuse until it is released.
*/ */
@Override @Override
public void onDelivered(String target) { public void onDelivered(String target, TurnToken token) {
MemberSession current = findByTerminal(target); MemberSession current = findByTerminal(target);
if (current == null) return; if (current == null) return;
if (current.state() != MemberSession.State.READY && current.state() != MemberSession.State.DONE) { if (current.state() != MemberSession.State.READY && current.state() != MemberSession.State.DONE) {