From e186c7945ab456ed10797e32d8a271463b32cf68 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 06:21:53 +0200 Subject: [PATCH 1/4] CB-577: correlate async questions to turns --- .../dev/ltms/bridged/msg/MessageService.java | 23 +++++++++++++------ .../ltms/bridged/msg/MessageServiceTest.java | 17 ++++++++++++++ 2 files changed, 33 insertions(+), 7 deletions(-) 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..cc99ff4 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -173,6 +173,9 @@ public final class MessageService { private final ConcurrentHashMap tasks = new ConcurrentHashMap<>(); /** Async tasks that have accepted delivery for a target. */ private final ConcurrentHashMap> asyncTasksByTarget = new ConcurrentHashMap<>(); + /** Async task that owns each exact forward rendezvous waiter. */ + private final ConcurrentHashMap, Task> asyncTasksByWaiter = + new ConcurrentHashMap<>(); /** Async tickets paused on a specific {@code bridge_ask} turn. */ private final ConcurrentHashMap asyncTasksByTurn = new ConcurrentHashMap<>(); private final AtomicLong ticketSeq = new AtomicLong(); @@ -385,6 +388,9 @@ public final class MessageService { // failed send leaves no stale waiter behind. CompletableFuture 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 +414,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 +440,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 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); @@ -613,17 +624,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 targetTasks = asyncTasksByTarget.get(target); - Task task = targetTasks == null ? null : targetTasks.stream() - .filter(candidate -> !candidate.future.isDone()) - .findFirst() - .orElse(null); + private Task markAsyncQuestion(CompletableFuture waiter, String text, String turnId) { + Task task = 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. */ @@ -634,6 +642,7 @@ public final class MessageService { if (forgetTurn) { asyncTasksByTurn.remove(turnId, task); task.turnId = null; + untrackAsyncTarget(task); } } } diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index ecb1b5d..5577eba 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -676,6 +676,23 @@ 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()); + } + private void assertFailedTicket(String ticket, String reason) throws Exception { MessageService.TaskView view = awaitTicketPhase(ticket, MessageService.Phase.FAILED); assertEquals(reason, view.detail()); From 5275922d1da971ace525d091723e0544f65e5ecd Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 06:22:22 +0200 Subject: [PATCH 2/4] CB-577: handle questions without async waiters --- bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 cc99ff4..024e274 100644 --- a/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java +++ b/bridged/src/main/java/dev/ltms/bridged/msg/MessageService.java @@ -625,7 +625,7 @@ public final class MessageService { /** Record the active question for an async ticket; blocking sends have no entry and stay unchanged. */ private Task markAsyncQuestion(CompletableFuture waiter, String text, String turnId) { - Task task = asyncTasksByWaiter.get(waiter); + Task task = waiter == null ? null : asyncTasksByWaiter.get(waiter); if (task != null) { task.question = new Reply(Outcome.QUESTION, text, turnId); task.turnId = turnId; From 927e0151d437a949c4d0790bc3819bef35e80a69 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 06:23:22 +0200 Subject: [PATCH 3/4] CB-577: test async question waiter ownership --- .../ltms/bridged/msg/MessageServiceTest.java | 35 +++++++++++++++++++ 1 file changed, 35 insertions(+) diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index 5577eba..c14001b 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -10,6 +10,9 @@ 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; @@ -693,6 +696,38 @@ class MessageServiceTest { assertEquals(MessageService.Phase.DONE, awaitTicketPhase(next, MessageService.Phase.DONE).phase()); } + @Test + @SuppressWarnings("unchecked") + void asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter() throws Exception { + String first = messages.sendAsync(T, "first task"); + awaitWaiting(); + String second = messages.sendAsync(T, "second task"); + + // Model the resolveQuestion/markAsyncQuestion race: another accepted task reached the target set. + Field tasksField = MessageService.class.getDeclaredField("tasks"); + tasksField.setAccessible(true); + Map tasks = (Map) tasksField.get(messages); + Field byTargetField = MessageService.class.getDeclaredField("asyncTasksByTarget"); + byTargetField.setAccessible(true); + Map> byTarget = (Map>) byTargetField.get(messages); + Set targetTasks = byTarget.get(T); + targetTasks.clear(); + targetTasks.add(tasks.get(second)); + + CompletableFuture ask = + CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000)); + assertEquals(MessageService.Phase.ASKING, awaitTicketPhase(first, MessageService.Phase.ASKING).phase()); + assertEquals(MessageService.Phase.PENDING, messages.poll(second).phase()); + + MessageService.TaskView asking = messages.poll(first); + CompletableFuture 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()); From 74b0087ebb54ceaa926a005a0f6af09f64945711 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Sat, 15 Aug 2026 06:24:05 +0200 Subject: [PATCH 4/4] CB-577: model async question ownership race --- .../java/dev/ltms/bridged/msg/MessageServiceTest.java | 10 ++++------ 1 file changed, 4 insertions(+), 6 deletions(-) diff --git a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java index c14001b..a72804e 100644 --- a/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java +++ b/bridged/src/test/java/dev/ltms/bridged/msg/MessageServiceTest.java @@ -701,23 +701,21 @@ class MessageServiceTest { void asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter() throws Exception { String first = messages.sendAsync(T, "first task"); awaitWaiting(); - String second = messages.sendAsync(T, "second task"); // Model the resolveQuestion/markAsyncQuestion race: another accepted task reached the target set. - Field tasksField = MessageService.class.getDeclaredField("tasks"); - tasksField.setAccessible(true); - Map tasks = (Map) tasksField.get(messages); Field byTargetField = MessageService.class.getDeclaredField("asyncTasksByTarget"); byTargetField.setAccessible(true); Map> byTarget = (Map>) byTargetField.get(messages); Set 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(tasks.get(second)); + targetTasks.add(constructor.newInstance(T)); CompletableFuture ask = CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000)); assertEquals(MessageService.Phase.ASKING, awaitTicketPhase(first, MessageService.Phase.ASKING).phase()); - assertEquals(MessageService.Phase.PENDING, messages.poll(second).phase()); MessageService.TaskView asking = messages.poll(first); CompletableFuture answer = CompletableFuture.supplyAsync(