ec3001796a
Record PrimaryRegistry delegator ownership via a MessageService accepted-delivery hook (won the session lock + queued delivery), never at bridge_send request time, so a concurrent sender that times out BUSY cannot steal a live turn's reply routing. Make Rendezvous.open atomic fail-if-present so a double open trips loudly instead of replacing the waiter another send is blocked on. Answering a bridge_ask keeps the same ownership (no rewrite). Adds ownership/rendezvous regression tests.
118 lines
5.0 KiB
Java
118 lines
5.0 KiB
Java
package dev.ltms.bridged.mcp;
|
|
|
|
import org.slf4j.Logger;
|
|
import org.slf4j.LoggerFactory;
|
|
|
|
import java.util.Optional;
|
|
import java.util.concurrent.ConcurrentHashMap;
|
|
import java.util.concurrent.atomic.AtomicReference;
|
|
|
|
/**
|
|
* Single-slot, thread-safe registry for the primary's herdr {@code terminal_id}.
|
|
*
|
|
* <p>Populated from the caller terminal of orchestration-side MCP tools
|
|
* ({@code bridge_send}, {@code bridge_spawn}) — tools that only the primary calls.
|
|
* A pinned terminal (from config) seeds the registry at construction and makes
|
|
* subsequent {@link #record(String)} calls no-ops.
|
|
*
|
|
* <p>The push loop ({@code ReplyPushLoop}) uses {@link #isKnown()} to decide
|
|
* whether active nudging is possible; an empty registry means the primary is
|
|
* off-host or non-herdr and delivery falls back to pull.
|
|
*/
|
|
public final class PrimaryRegistry {
|
|
|
|
private static final Logger log = LoggerFactory.getLogger(PrimaryRegistry.class);
|
|
|
|
private final AtomicReference<String> terminal = new AtomicReference<>();
|
|
private final boolean pinned;
|
|
|
|
/**
|
|
* CB-532: worker terminal → the lead that delegated to it. The single slot above answers "who is
|
|
* THE primary", a question with no correct answer once two leads orchestrate the same fleet:
|
|
* whichever called {@code bridge_send} first captured every nudge, including nudges for the
|
|
* other lead's delegations. This map answers the question that actually matters — "who is
|
|
* waiting on THIS worker" — and is what lets {@code primary.terminal} be retired.
|
|
*/
|
|
private final ConcurrentHashMap<String, String> leadByTarget = new ConcurrentHashMap<>();
|
|
|
|
/**
|
|
* @param pinnedTerminal an optional pinned terminal from config ({@code null}/blank = unpinned)
|
|
*/
|
|
public PrimaryRegistry(String pinnedTerminal) {
|
|
if (pinnedTerminal != null && !pinnedTerminal.isBlank()) {
|
|
this.terminal.set(pinnedTerminal);
|
|
this.pinned = true;
|
|
log.info("primary terminal pinned: {}", pinnedTerminal);
|
|
} else {
|
|
this.pinned = false;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Record a terminal_id. No-op when:
|
|
* <ul>
|
|
* <li>the registry is pinned (config override),
|
|
* <li>{@code terminalId} is {@code null} or blank (non-herdr caller).
|
|
* </ul>
|
|
*/
|
|
public void record(String terminalId) {
|
|
if (pinned) return;
|
|
if (terminalId == null || terminalId.isBlank()) return;
|
|
String prev = terminal.getAndSet(terminalId);
|
|
if (prev == null) {
|
|
log.debug("primary terminal learned: {}", terminalId);
|
|
} else if (!prev.equals(terminalId)) {
|
|
log.debug("primary terminal changed: {} -> {}", prev, terminalId);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Record that {@code leadTerminal} owns the accepted delegation of worker {@code target} (CB-532).
|
|
*
|
|
* <p>Called from the {@code MessageService} accepted-delivery hook — only after a send has won
|
|
* the session's send lock and queued delivery — where both halves are known (CB-548). It is
|
|
* deliberately <em>not</em> called at {@code bridge_send} request time: a concurrent sender that
|
|
* times out {@code BUSY} must not steal a live delegation's reply routing without ever owning
|
|
* the turn. Last writer wins — if a second lead's later send is accepted, replies follow the
|
|
* lead that most recently delegated to it, which is the one waiting.
|
|
*/
|
|
public void recordDelegation(String target, String leadTerminal) {
|
|
if (target == null || target.isBlank() || leadTerminal == null || leadTerminal.isBlank()) {
|
|
return;
|
|
}
|
|
leadByTarget.put(target, leadTerminal);
|
|
}
|
|
|
|
/** Forget a worker's delegating lead — call on release, so a torn-down session leaks nothing. */
|
|
public void forgetDelegation(String target) {
|
|
if (target != null) {
|
|
leadByTarget.remove(target);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Where a nudge about {@code target}'s reply should go: the lead that delegated to it, falling
|
|
* back to the single known primary.
|
|
*
|
|
* <p>The fallback matters after a daemon restart, which loses the map while the durable inbox
|
|
* keeps the reply. With one lead the fallback is unambiguous and correct. With several and no
|
|
* recorded delegation there is no right answer, so this returns empty rather than guessing —
|
|
* delivery degrades to pull, which is exactly what the durable inbox is for, instead of
|
|
* interrupting the wrong lead with someone else's result.
|
|
*/
|
|
public Optional<String> nudgeTargetFor(String target) {
|
|
String lead = target == null ? null : leadByTarget.get(target);
|
|
return lead != null ? Optional.of(lead) : Optional.ofNullable(terminal.get());
|
|
}
|
|
|
|
/** The known primary terminal, or empty if not yet learned (and not pinned). */
|
|
public Optional<String> primaryTerminal() {
|
|
return Optional.ofNullable(terminal.get());
|
|
}
|
|
|
|
/** {@code true} once a terminal has been recorded (or was pinned at construction). */
|
|
public boolean isKnown() {
|
|
return terminal.get() != null;
|
|
}
|
|
}
|