From 38f13f4b3c8289ee2cb88dca3ba5c10fc40115f8 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 06:27:39 +0200 Subject: [PATCH] CB-573: carry accepted turn tokens on delivery --- .../main/java/dev/ltms/bridged/Bridged.java | 6 +++--- .../bridged/inject/CompletionResolver.java | 9 +++++---- .../dev/ltms/bridged/inject/Injector.java | 9 +++++---- .../dev/ltms/bridged/inject/TurnListener.java | 4 +++- .../dev/ltms/bridged/msg/MessageService.java | 3 ++- .../java/dev/ltms/bridged/msg/TurnToken.java | 20 +++++++++++++++++++ .../ltms/bridged/session/SessionManager.java | 3 ++- 7 files changed, 40 insertions(+), 14 deletions(-) create mode 100644 bridged/src/main/java/dev/ltms/bridged/msg/TurnToken.java diff --git a/bridged/src/main/java/dev/ltms/bridged/Bridged.java b/bridged/src/main/java/dev/ltms/bridged/Bridged.java index fbae9b7..70a10d9 100644 --- a/bridged/src/main/java/dev/ltms/bridged/Bridged.java +++ b/bridged/src/main/java/dev/ltms/bridged/Bridged.java @@ -276,9 +276,9 @@ public final class Bridged { } @Override - public void onDelivered(String target) { - completion.onDelivered(target); - sessions.onDelivered(target); + public void onDelivered(String target, dev.ltms.bridged.msg.TurnToken token) { + completion.onDelivered(target, token); + sessions.onDelivered(target, token); } @Override diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java index e504b96..5259ffd 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/CompletionResolver.java @@ -2,6 +2,7 @@ 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; @@ -83,17 +84,17 @@ public final class CompletionResolver implements TurnListener { } @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 // 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); + captureBaseline(target, token); } /** Capture the in-flight turn: its waiter and pre-turn baseline (the testable core of {@link #onDelivered}). */ - void captureBaseline(String target) { - CompletableFuture waiter = rendezvous.currentWaiter(target); + void captureBaseline(String target, TurnToken token) { + CompletableFuture waiter = token.waiter(); if (waiter == null) { inFlight.remove(target); // no send is waiting on this delivery — nothing to resolve later return; diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java index 6144c6b..eeb62f3 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/Injector.java @@ -2,6 +2,7 @@ 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; @@ -129,7 +130,7 @@ public final class Injector { } /** A pending message and the future that completes when it has been delivered. */ - private record Pending(String text, CompletableFuture delivered) { + private record Pending(String text, TurnToken token, CompletableFuture delivered) { } /** Per-worker delivery state, guarded by its own monitor (single writer per worker). */ @@ -159,9 +160,9 @@ public final class Injector { *

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 enqueue(String target, String text) { + public CompletableFuture enqueue(String target, String text, TurnToken token) { CompletableFuture delivered = new CompletableFuture<>(); - Pending p = new Pending(text, delivered); + Pending p = new Pending(text, token, 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 @@ -356,7 +357,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); + turnListener.onDelivered(target, sent.token()); sent.delivered().complete(null); } } diff --git a/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java b/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java index 68e377d..7d453ee 100644 --- a/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java +++ b/bridged/src/main/java/dev/ltms/bridged/inject/TurnListener.java @@ -1,5 +1,7 @@ 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 @@ -57,7 +59,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) { + default void onDelivered(String target, TurnToken token) { } /** No-op default for callers that only need delivery, not completion signalling. */ 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 c50675e..13fa40d 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -385,13 +385,14 @@ public final class MessageService { // failed send leaves no stale waiter behind. CompletableFuture reply = rendezvous.open(target); try { + TurnToken token = new TurnToken(target, reply); // 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 delivered = injector.enqueue(target, content); + CompletableFuture delivered = injector.enqueue(target, content, token); try { Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS); return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId())); diff --git a/bridged/src/main/java/dev/ltms/bridged/msg/TurnToken.java b/bridged/src/main/java/dev/ltms/bridged/msg/TurnToken.java new file mode 100644 index 0000000..f3f8c2e --- /dev/null +++ b/bridged/src/main/java/dev/ltms/bridged/msg/TurnToken.java @@ -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 waiter; + + public TurnToken(String target, CompletableFuture waiter) { + this.target = target; + this.waiter = waiter; + } + + public String target() { return target; } + public CompletableFuture waiter() { return waiter; } +} diff --git a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java index fdb5766..9ebc167 100644 --- a/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java +++ b/bridged/src/main/java/dev/ltms/bridged/session/SessionManager.java @@ -4,6 +4,7 @@ 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; @@ -425,7 +426,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) { + public void onDelivered(String target, TurnToken token) { MemberSession current = findByTerminal(target); if (current == null) return; if (current.state() != MemberSession.State.READY && current.state() != MemberSession.State.DONE) { -- 2.52.0