CB-520: split ReplyInbox into explicit own/release and publish halves

This commit is contained in:
Dai Ha
2026-08-04 17:44:21 +02:00
parent c1173346ef
commit 049ce4828c
11 changed files with 270 additions and 60 deletions
@@ -211,11 +211,16 @@ public final class Bridged {
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
pushLoop, metrics);
// CB-520: the reply inbox only consumes for agents this gateway owns. own on acquire,
// release on teardown. Do this before CB-516 so the inbox is owned before any reply can land.
sessions.onAcquire(replyInbox::own);
// CB-516: releasing a worker must fail whatever send was waiting on it. Without this a
// torn-down delegation kept reporting PENDING until the 30-minute async timeout, and never
// reached /metrics — the delegation was unresolvable and nothing said so.
sessions.onRelease(terminal ->
messages.abandon(terminal, "the worker session was released before it replied"));
sessions.onRelease(terminal -> {
messages.abandon(terminal, "the worker session was released before it replied");
replyInbox.release(terminal);
});
// MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp.
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
@@ -14,25 +14,30 @@ import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
/**
* AMQP-backed {@link ReplyInbox} (CB-307 Stage 2): genuine cross-restart durability behind the same
* port {@link InMemoryReplyInbox} implements as soft state.
*
* <p><strong>Mapping — consume-and-hold with deferred manual ack.</strong> Each target owns a durable
* queue {@code agent.<target>.inbox}. A manual-ack consumer pulls persistent messages off that queue
* into an in-memory <em>held</em> map (keyed by {@code msgId}) but does <em>not</em> ack them.
* {@link #peek} returns that snapshot; {@link #ack} acks the broker delivery-tag and drops the entry.
* Because messages stay unacked until the primary actually drains them, a crash (or a {@code java -jar}
* bounce) before caller-ack leaves them on the broker — it redelivers on reconnect. That is the
* durability the in-memory adapter cannot give, with the port contract preserved.
* <p><strong>Mapping — consume-and-hold with deferred manual ack.</strong> Each target has a durable
* queue {@code agent.<target>.inbox}. The gateway that owns the target starts a manual-ack consumer
* ({@link #own}) that pulls persistent messages off that queue into an in-memory <em>held</em> map
* (keyed by {@code msgId}) but does <em>not</em> ack them. {@link #peek} returns that snapshot;
* {@link #ack} acks the broker delivery-tag and drops the entry. Because messages stay unacked until
* the owning gateway actually drains them, a crash (or a {@code java -jar} bounce) before caller-ack
* leaves them on the broker — it redelivers on reconnect. That is the durability the in-memory
* adapter cannot give, with the port contract preserved.
*
* <p><strong>Ownership is explicit.</strong> {@link #own} declares the queue and starts the consumer;
* {@link #release} cancels it. {@link #publish} sends to the queue but does <em>not</em> imply ownership
* and does not attach a consumer. This split is required by CB-308 federation, where one gateway may
* publish to an agent owned by another gateway; in that case the publisher must not compete for
* deliveries.
*
* <p><strong>Dedup.</strong> The consumer keys the held map by {@code msgId}; a redelivered duplicate
* (at-least-once, or a producer double-publish) is acked-and-dropped on arrival, so it never
* double-queues. {@link #publish} additionally short-circuits an already-held {@code msgId} — a
* fast path; the consumer-side check is the real guarantee.
* double-queues.
*
* <p><strong>Visibility.</strong> Unlike the in-memory adapter, publish → broker → consumer is
* asynchronous, so a {@link #peek} immediately after {@link #publish} may not yet see the message
@@ -51,12 +56,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
private final Connection connection;
private final Channel channel;
/** All channel operations (publish/declare/ack) serialize on this — a Channel is not thread-safe. */
/** All channel operations (publish/declare/ack/cancel) serialize on this — a Channel is not thread-safe. */
private final Object channelLock = new Object();
/** target → (msgId → held delivery). Per-target map is guarded by synchronizing on itself. */
private final ConcurrentHashMap<String, LinkedHashMap<String, Held>> held = new ConcurrentHashMap<>();
/** Targets whose queue is declared and consumer is running. */
private final Set<String> consuming = ConcurrentHashMap.newKeySet();
/** Targets whose queue is declared and consumer is running, mapped to their broker consumer tag. */
private final ConcurrentHashMap<String, String> consumerTags = new ConcurrentHashMap<>();
/** A message pulled off the broker but not yet acked: its delivery-tag plus the port payload. */
private record Held(long deliveryTag, InboxMessage message) {}
@@ -103,16 +108,41 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
@Override
public void publish(String target, String msgId, String content) {
ensureConsuming(target);
var perTarget = held.get(target);
if (perTarget != null) {
synchronized (perTarget) {
if (perTarget.containsKey(msgId)) {
return; // already held — producer-side fast dedup
}
public void own(String target) {
synchronized (channelLock) {
if (consumerTags.containsKey(target)) {
return; // already owning this target
}
String queue = queueName(target);
try {
channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
String tag = channel.basicConsume(queue, false, deliverCallback(target), _ -> { });
consumerTags.put(target, tag);
log.debug("AMQP inbox owns queue {} for target {}", queue, target);
} catch (IOException e) {
throw new IllegalStateException("cannot own queue " + queue, e);
}
}
}
@Override
public void release(String target) {
synchronized (channelLock) {
String tag = consumerTags.remove(target);
held.remove(target); // stale delivery tags must not survive release
if (tag == null) {
return;
}
try {
channel.basicCancel(tag);
} catch (IOException e) {
throw new IllegalStateException("cannot cancel consumer for " + target, e);
}
}
}
@Override
public void publish(String target, String msgId, String content) {
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.messageId(msgId)
.deliveryMode(2) // persistent — survives a broker restart
@@ -129,7 +159,6 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
@Override
public List<InboxMessage> peek(String target) {
ensureConsuming(target);
var perTarget = held.get(target);
if (perTarget == null) {
return List.of();
@@ -166,26 +195,6 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
}
/** Declare the durable per-target queue and start its manual-ack consumer, once per target. */
private void ensureConsuming(String target) {
if (consuming.contains(target)) {
return;
}
synchronized (channelLock) {
if (!consuming.add(target)) {
return; // another thread just set it up
}
String queue = queueName(target);
try {
channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
channel.basicConsume(queue, false, deliverCallback(target), _ -> { });
} catch (IOException e) {
consuming.remove(target);
throw new IllegalStateException("cannot consume queue " + queue, e);
}
}
}
private DeliverCallback deliverCallback(String target) {
return (_, delivery) -> {
String msgId = delivery.getProperties().getMessageId();
@@ -2,6 +2,7 @@ package dev.ltms.bridged.msg;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
/**
@@ -9,12 +10,29 @@ import java.util.concurrent.ConcurrentHashMap;
* Per-target FIFO ordering (insertion order via {@link LinkedHashMap}). Dedup by {@code msgId}
* within a target. Thread-safe for concurrent publish vs. drain.
*
* <p><strong>Ownership is explicit.</strong> {@link #own} marks a target as locally owned so that
* {@link #peek} and {@link #ack} operate on it; {@link #publish} works whether or not the target is
* owned. {@link #release} clears the local snapshot. This mirrors the AMQP adapter's contract so the
* non-broker path stays interchangeable.
*
* <p><strong>This is soft-state, NOT persistence.</strong> Lost on a {@code java -jar} bounce — that
* is correct and consistent with "bridged stays soft-state." The Stage-2 AMQP adapter replaces this.
*/
public final class InMemoryReplyInbox implements ReplyInbox {
private final ConcurrentHashMap<String, LinkedHashMap<String, InboxMessage>> store = new ConcurrentHashMap<>();
private final Set<String> owned = ConcurrentHashMap.newKeySet();
@Override
public void own(String target) {
owned.add(target);
}
@Override
public void release(String target) {
owned.remove(target);
store.remove(target);
}
@Override
public void publish(String target, String msgId, String content) {
@@ -27,6 +45,9 @@ public final class InMemoryReplyInbox implements ReplyInbox {
@Override
public List<InboxMessage> peek(String target) {
if (!owned.contains(target)) {
return List.of();
}
var perTarget = store.get(target);
if (perTarget == null) {
return List.of();
@@ -39,6 +60,9 @@ public final class InMemoryReplyInbox implements ReplyInbox {
@Override
public void ack(String target, String msgId) {
if (!owned.contains(target)) {
return;
}
var perTarget = store.get(target);
if (perTarget != null) {
//noinspection SynchronizationOnLocalVariableOrMethodParameter
@@ -10,16 +10,36 @@ import java.util.List;
* <p><strong>This interface is the port.</strong> {@link InMemoryReplyInbox} is the Stage-1 adapter;
* an AMQP-backed adapter (Stage 2) must implement the same contract (idempotent publish, FIFO peek,
* at-least-once ack).
*
* <p><strong>Ownership is explicit.</strong> A gateway {@link #own owns} the inbox for each agent it
* spawned; only the owner consumes and drains it. {@link #publish} sends a reply to the target's
* inbox but does <em>not</em> imply ownership or start a consumer. This separation is required by
* CB-308 federation, where one gateway may publish to an agent owned by another gateway.
*/
public interface ReplyInbox {
/** A queued reply: an idempotency id, the worker session it came from, and the reply text. */
record InboxMessage(String msgId, String target, String content) {}
/**
* Start owning (consuming) the inbox for {@code target}. Idempotent: multiple calls for the same
* target are no-ops. The owner is the only gateway that may {@link #peek} and {@link #ack} replies
* for this target.
*/
void own(String target);
/**
* Stop owning (consuming) the inbox for {@code target}. Idempotent. Any replies held locally but
* not yet acked are dropped from the local snapshot; the underlying durable queue keeps
* unacked messages for redelivery when the target is re-owned.
*/
void release(String target);
/**
* Queue {@code content} from worker {@code target} under {@code msgId}. Idempotent: publishing an
* already-present {@code msgId} for {@code target} is a no-op (dedup), so an at-least-once Stage-2
* redelivery cannot double-queue.
* redelivery cannot double-queue. Publishing does <em>not</em> imply ownership and must not start a
* consumer.
*/
void publish(String target, String msgId, String content);
@@ -48,8 +48,10 @@ public final class SessionManager implements TurnListener {
private final LongSupplier nowNanos;
private final int contextCap;
/** CB-520: notified with a terminalId on every acquire; no-op until wired. */
private final List<Consumer<String>> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
/** CB-516: notified with a terminalId on every release; no-op until wired. */
private volatile Consumer<String> releaseListener = _ -> { };
private final List<Consumer<String>> releaseListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
/** Backward-compatible constructor: shared-tree sessions, production git seam. */
public SessionManager(PeerLauncher launcher) {
@@ -129,6 +131,7 @@ public final class SessionManager implements TurnListener {
registry.put(handle.id(), session);
log.debug("acquired session id={} terminal={} profile={} owner={}",
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
notifyAcquired(session.terminalId());
return session;
}
return acquireWithWorktree(profile, requestedCwd, callerCwd, ownerTerminal, wt);
@@ -151,18 +154,44 @@ public final class SessionManager implements TurnListener {
}
}
/**
* Register a callback invoked with a session's {@code terminalId} whenever it is acquired
* (CB-520). This is the hook that lets the reply inbox {@code own} a target's queue.
*/
public void onAcquire(Consumer<String> listener) {
if (listener != null) {
acquireListeners.add(listener);
}
}
/**
* Register a callback invoked with a session's {@code terminalId} whenever it is released
* (CB-516). Every teardown path funnels through {@link #release}, so one hook covers the REST
* and MCP stop tools, the idle-TTL reaper, {@code recycle}, and shutdown drain alike.
*
* <p>Set rather than injected because {@code MessageService} — the intended listener — is
* <p>Added rather than injected because {@code MessageService} — one intended listener — is
* constructed after this manager (it needs the injector and rendezvous, which need the session
* presence view this manager exposes). Wiring it at construction would require breaking that
* cycle for one callback.
*/
public void onRelease(Consumer<String> listener) {
this.releaseListener = (listener == null) ? _ -> { } : listener;
if (listener != null) {
releaseListeners.add(listener);
}
}
/** A listener failure must never prevent the acquisition it is reacting to. */
private void notifyAcquired(String terminalId) {
if (terminalId == null) {
return;
}
for (Consumer<String> listener : acquireListeners) {
try {
listener.accept(terminalId);
} catch (RuntimeException e) {
log.warn("acquire listener failed for terminal {}: {}", terminalId, e.toString());
}
}
}
/** A listener failure must never prevent the teardown it is reacting to. */
@@ -170,10 +199,12 @@ public final class SessionManager implements TurnListener {
if (terminalId == null) {
return;
}
try {
releaseListener.accept(terminalId);
} catch (RuntimeException e) {
log.warn("release listener failed for terminal {}: {}", terminalId, e.toString());
for (Consumer<String> listener : releaseListeners) {
try {
listener.accept(terminalId);
} catch (RuntimeException e) {
log.warn("release listener failed for terminal {}: {}", terminalId, e.toString());
}
}
}
@@ -221,6 +252,7 @@ public final class SessionManager implements TurnListener {
registry.put(handle.id(), session);
log.debug("acquired worktree session id={} terminal={} profile={} branch={} path={}",
handle.id(), handle.terminalId(), session.profile(), session.branch(), session.worktree());
notifyAcquired(session.terminalId());
return session;
}