Compare commits
17 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| bb750cdba3 | |||
| d56c77b368 | |||
| 83e2ff06cf | |||
| f5deaafd06 | |||
| fe46311266 | |||
| 08968bb1b7 | |||
| 27bbd11f06 | |||
| 48d7841fbf | |||
| cc1df11f69 | |||
| fdfd4ac491 | |||
| d5dd5639ae | |||
| 7a120b3256 | |||
| 863d477966 | |||
| cec48832be | |||
| a36b7ccd7c | |||
| 16de9df000 | |||
| 81a0cf4710 |
@@ -105,8 +105,8 @@ the merge — and merging on a reviewer's word is delegating it by proxy.
|
||||
|---|---|
|
||||
| Confirm your own role | `bridge_whoami` |
|
||||
| See backends available | `bridge_profiles` |
|
||||
| Start a member | `bridge_spawn{role?, profile?, cwd?, worktree?, ticket?}` → `sessionId` + `paneId` |
|
||||
| See the fleet | `bridge_list` → `leads` (your peers) + `members` · one peer's state: `bridge_status{sessionId}` |
|
||||
| Start a member | `bridge_spawn{role?, profile?, cwd?, worktree?, ticket?, sessionName?, resumeSessionId?}` → `sessionId` + `paneId` |
|
||||
| See the fleet | `bridge_list` → `leads` (your peers) + `members` (each carries `agentSessionId` when its backend knows one) · one peer's state: `bridge_status{sessionId}` |
|
||||
| Delegate (blocking) | `bridge_send{sessionId, content}` |
|
||||
| Delegate (long task) | `bridge_send{sessionId, content, wait:false}` → ticket → `bridge_poll{ticket}` |
|
||||
| Answer a member's `bridge_ask` | `bridge_send{turnId, content}` — **not** `sessionId` |
|
||||
|
||||
@@ -275,6 +275,26 @@ profiles:
|
||||
# argv: ["opencode"]
|
||||
# How an unqualified spawn chooses a profile: fixed (default, reproduces pre-CB-518 behaviour),
|
||||
# round-robin, or weighted. Omitting this key is a strict no-op for existing configs.
|
||||
#
|
||||
# `weighted` IS NOT "cheapest first" — read this before you set weights (CB-589).
|
||||
# It is smooth weighted round-robin: it spreads spawns across EVERY profile that has a free slot,
|
||||
# in weight ratio. It has no idea which profile costs money. So with local:10 / paid:2 you do not
|
||||
# get "use local, overflow to paid" — you get roughly one spawn in six going to the paid profile
|
||||
# while the local box still has a free slot.
|
||||
#
|
||||
# There is a sharper second effect. The policy's running score map lives for the daemon's whole
|
||||
# life. While a profile is at maxLoad it is filtered out and its score FREEZES, so the paid
|
||||
# profiles keep accumulating against it. When the local slot frees up it returns with a stale
|
||||
# score and can LOSE the next pick — a paid spawn while the free box sits idle.
|
||||
#
|
||||
# Until a real cost-first policy exists, the workaround is to make the ratio decisive rather than
|
||||
# proportional: give the free profile a weight so large that it wins every pick it is eligible
|
||||
# for, and paid profiles only ever take genuine overflow. On this host that is local weight 100
|
||||
# against paid weights of ~1.
|
||||
#
|
||||
# The gotcha with that workaround: it expresses a PREFERENCE ORDER through a RATIO knob. Add a
|
||||
# future profile at weight 150 and it silently outranks the free box, with nothing to warn you.
|
||||
# Re-check the weights whenever you add a profile.
|
||||
placement: weighted
|
||||
|
||||
# How long a credential sits out after a BACKEND_EXHAUSTED classification (CB-578 stage B), in
|
||||
|
||||
@@ -441,6 +441,11 @@ public final class Bridged {
|
||||
reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch()
|
||||
+ " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none");
|
||||
}
|
||||
// CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the
|
||||
// member's conversation instead of only re-dispatching a fresh one onto the same files.
|
||||
if (detail.agentSessionId() != null) {
|
||||
reason += " agentSessionId=" + detail.agentSessionId();
|
||||
}
|
||||
messages.abandon(detail.terminalId(), reason);
|
||||
replyInbox.release(detail.terminalId());
|
||||
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
|
||||
@@ -555,7 +560,10 @@ public final class Bridged {
|
||||
* where set (opt-in). Derived from the config, not hard-coded, so a new profile is covered for
|
||||
* free. A var required by more than one profile is one entry naming every profile that needs
|
||||
* it. Deliberately excludes {@code auth.tokenEnv}: that one is already enforced loudly, by a
|
||||
* startup throw, a few lines above this method's call site.
|
||||
* startup throw in {@code main()} — about 370 lines <em>below</em> this method's call site
|
||||
* ({@link #reportRequiredSecrets(BridgedConfig)}), not a few lines above it. That throw only
|
||||
* fires when {@code auth.mode: token} is configured; under the default loopback-trust mode it
|
||||
* never runs, and {@code auth.tokenEnv} is simply not required.
|
||||
*
|
||||
* <p>Package-private and pure (no I/O, no logging) so the derivation is unit-testable without
|
||||
* capturing log output; {@link #reportRequiredSecrets(BridgedConfig)} is the logging caller.
|
||||
|
||||
@@ -7,6 +7,7 @@ import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
|
||||
import dev.ltms.bridged.msg.AmqpReplyInbox;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
import dev.ltms.bridged.placement.PlacementPolicies;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -986,7 +987,14 @@ public record BridgedConfig(
|
||||
warnUnknownTopLevelKeys(yaml, path);
|
||||
rejectDuplicateMemberSlots(yaml);
|
||||
rejectNegativeMaxLoad(yaml);
|
||||
rejectUnknownKind(yaml);
|
||||
rejectUnknownAuthMode(yaml);
|
||||
rejectUnknownPlacement(yaml);
|
||||
BridgedConfig cfg = YAML.readValue(yaml, BridgedConfig.class);
|
||||
// CB-606: validated here, eagerly, using PlacementPolicies.fromName as the single source
|
||||
// of truth — not lazily at first spawn (see CompositePeerLauncher's placementPolicy
|
||||
// Supplier), where a bad name would still start a daemon that looks healthy.
|
||||
rejectUnknownPlacementPolicy(cfg.placement());
|
||||
return cfg.withDefaults();
|
||||
} catch (IOException e) {
|
||||
throw new UncheckedIOException("cannot read bridged config at " + path, e);
|
||||
@@ -1280,6 +1288,162 @@ public record BridgedConfig(
|
||||
}
|
||||
}
|
||||
|
||||
/** The peer kinds this build has an adapter for — {@link Profile#kind()}'s only valid values. */
|
||||
private static final Set<String> KNOWN_KINDS = Set.of(Profile.KIND_CLAUDE_CODE, Profile.KIND_OPENCODE);
|
||||
|
||||
/**
|
||||
* Reject a profile whose {@code kind:} is not one of {@link #KNOWN_KINDS} (CB-604), naming the
|
||||
* profile, the value it set, and the accepted set.
|
||||
*
|
||||
* <p>{@link Profile}'s compact constructor only lower-cases {@code kind} and compares it against
|
||||
* {@code KIND_OPENCODE} — anything else, including a typo like {@code opencod}, silently falls
|
||||
* into the claude-code bucket ({@link dev.ltms.bridged.member.CompositePeerLauncher} routes by
|
||||
* exact adapter claim, not by membership in a known set). With {@code argv:} also unset, the argv
|
||||
* default special-cases only the exact string {@code "claude-code"}, so the launch command falls
|
||||
* back to {@code List.of(kind)} — the daemon then tries to run a program literally named after the
|
||||
* typo. {@code CompositePeerLauncher}'s constructor already treats a profile claimed by two
|
||||
* adapters as fatal (CB-402); an unrecognized kind is the same class of adapter-routing mistake
|
||||
* and gets the same treatment here, at config load, rather than surfacing later as a failed spawn.
|
||||
*
|
||||
* @param yaml the raw config text
|
||||
* @throws IllegalStateException when any profile's {@code kind} is a non-blank value not in
|
||||
* {@link #KNOWN_KINDS} (case-insensitive)
|
||||
*/
|
||||
static void rejectUnknownKind(String yaml) {
|
||||
Map<?, ?> raw;
|
||||
try {
|
||||
raw = YAML.readValue(yaml, Map.class);
|
||||
} catch (IOException | IllegalArgumentException e) {
|
||||
return; // a malformed file is reported by the real parse, not here
|
||||
}
|
||||
if (raw == null || !(raw.get("profiles") instanceof Map<?, ?> profiles)) {
|
||||
return;
|
||||
}
|
||||
List<String> bad = profiles.entrySet().stream()
|
||||
.filter(e -> e.getValue() instanceof Map<?, ?> p
|
||||
&& p.get("kind") instanceof String k && !k.isBlank()
|
||||
&& !KNOWN_KINDS.contains(k.toLowerCase()))
|
||||
.map(e -> String.valueOf(e.getKey()) + "=" + ((Map<?, ?>) e.getValue()).get("kind"))
|
||||
.sorted()
|
||||
.toList();
|
||||
if (!bad.isEmpty()) {
|
||||
throw new IllegalStateException("refusing to start: profile(s) [" + String.join(", ", bad)
|
||||
+ "] set an unrecognized kind — accepted values are "
|
||||
+ String.join(", ", KNOWN_KINDS.stream().sorted().toList())
|
||||
+ " (case-insensitive); an unrecognized kind would otherwise fall back to the"
|
||||
+ " claude-code adapter and try to launch a program named after the typo.");
|
||||
}
|
||||
}
|
||||
|
||||
/** The auth modes this build understands — {@link Auth#mode()}'s only valid values. */
|
||||
private static final Set<String> KNOWN_AUTH_MODES = Set.of(Auth.MODE_LOOPBACK_TRUST, Auth.MODE_TOKEN);
|
||||
|
||||
/**
|
||||
* Reject an {@code auth.mode} that is not one of {@link #KNOWN_AUTH_MODES} (CB-606), naming the
|
||||
* value and the accepted set.
|
||||
*
|
||||
* <p>{@link Auth}'s compact constructor only lower-cases {@code mode}, and
|
||||
* {@link Auth#tokenMode()} only compares the result against {@code MODE_TOKEN} — anything else,
|
||||
* including a typo like {@code toekn}, silently behaves as {@code loopback-trust}. That fallback
|
||||
* is otherwise checked only by {@link #validateAuthExposure()}, and only when the bind is
|
||||
* non-loopback: on a loopback bind (the common case) the typo is invisible end to end — the
|
||||
* daemon starts cleanly and authenticates nobody while the operator believes {@code token} mode
|
||||
* is active. Refuse it here, unconditionally, at config load, rather than let it hide behind the
|
||||
* bind check.
|
||||
*
|
||||
* @param yaml the raw config text
|
||||
* @throws IllegalStateException when {@code auth.mode} is a non-blank value not in
|
||||
* {@link #KNOWN_AUTH_MODES} (case-insensitive)
|
||||
*/
|
||||
static void rejectUnknownAuthMode(String yaml) {
|
||||
Map<?, ?> raw;
|
||||
try {
|
||||
raw = YAML.readValue(yaml, Map.class);
|
||||
} catch (IOException | IllegalArgumentException e) {
|
||||
return; // a malformed file is reported by the real parse, not here
|
||||
}
|
||||
if (raw == null || !(raw.get("auth") instanceof Map<?, ?> auth)) {
|
||||
return;
|
||||
}
|
||||
if (!(auth.get("mode") instanceof String mode) || mode.isBlank()
|
||||
|| KNOWN_AUTH_MODES.contains(mode.toLowerCase())) {
|
||||
return;
|
||||
}
|
||||
throw new IllegalStateException("refusing to start: auth.mode=" + mode
|
||||
+ " is not recognized — accepted values are "
|
||||
+ String.join(", ", KNOWN_AUTH_MODES.stream().sorted().toList())
|
||||
+ " (case-insensitive); an unrecognized mode would otherwise silently fall back to"
|
||||
+ " loopback-trust, which authenticates nobody.");
|
||||
}
|
||||
|
||||
/** The per-profile placements this build understands — {@link Profile#placement()}'s only valid values. */
|
||||
private static final Set<String> KNOWN_PLACEMENTS = Set.of("tab", "pane");
|
||||
|
||||
/**
|
||||
* Reject a profile whose {@code placement:} is not one of {@link #KNOWN_PLACEMENTS} (CB-606),
|
||||
* naming the profile, the value it set, and the accepted set.
|
||||
*
|
||||
* <p>{@link Profile}'s compact constructor only lower-cases {@code placement}, and
|
||||
* {@link Profile#tabPlacement()} only compares the result against {@code "tab"} — anything else,
|
||||
* including a typo like {@code tabb}, silently falls back to the legacy pane placement with no
|
||||
* signal anywhere.
|
||||
*
|
||||
* @param yaml the raw config text
|
||||
* @throws IllegalStateException when any profile's {@code placement} is a non-blank value not in
|
||||
* {@link #KNOWN_PLACEMENTS} (case-insensitive)
|
||||
*/
|
||||
static void rejectUnknownPlacement(String yaml) {
|
||||
Map<?, ?> raw;
|
||||
try {
|
||||
raw = YAML.readValue(yaml, Map.class);
|
||||
} catch (IOException | IllegalArgumentException e) {
|
||||
return; // a malformed file is reported by the real parse, not here
|
||||
}
|
||||
if (raw == null || !(raw.get("profiles") instanceof Map<?, ?> profiles)) {
|
||||
return;
|
||||
}
|
||||
List<String> bad = profiles.entrySet().stream()
|
||||
.filter(e -> e.getValue() instanceof Map<?, ?> p
|
||||
&& p.get("placement") instanceof String pl && !pl.isBlank()
|
||||
&& !KNOWN_PLACEMENTS.contains(pl.toLowerCase()))
|
||||
.map(e -> String.valueOf(e.getKey()) + "=" + ((Map<?, ?>) e.getValue()).get("placement"))
|
||||
.sorted()
|
||||
.toList();
|
||||
if (!bad.isEmpty()) {
|
||||
throw new IllegalStateException("refusing to start: profile(s) [" + String.join(", ", bad)
|
||||
+ "] set an unrecognized placement — accepted values are "
|
||||
+ String.join(", ", KNOWN_PLACEMENTS.stream().sorted().toList())
|
||||
+ " (case-insensitive); an unrecognized placement would otherwise fall back to"
|
||||
+ " legacy pane placement with no signal anywhere.");
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Reject a top-level {@code placement:} policy name {@link PlacementPolicies#fromName} does not
|
||||
* recognize (CB-606), at config load rather than lazily at first spawn.
|
||||
*
|
||||
* <p>{@code CompositePeerLauncher} only calls {@link PlacementPolicies#fromName} per spawn,
|
||||
* through a {@code Supplier} that re-reads live config (CB-559, so a hot-reloaded placement
|
||||
* policy takes effect without a restart) — so a bad name still starts a daemon that looks
|
||||
* healthy and fails only the first time something spawns without naming a profile. Every other
|
||||
* field this class validates fails here, at load; this one gets the same treatment, calling
|
||||
* {@link PlacementPolicies#fromName} itself as the single source of truth for what is valid
|
||||
* rather than duplicating its accepted set.
|
||||
*
|
||||
* @param placement the raw, possibly null/blank {@code placement} value as parsed (before
|
||||
* {@link #withDefaults()} runs); {@code fromName} itself treats null/blank as
|
||||
* {@code fixed}, so this call changes no default
|
||||
* @throws IllegalStateException when {@code placement} is a name {@link PlacementPolicies} does
|
||||
* not recognize
|
||||
*/
|
||||
private static void rejectUnknownPlacementPolicy(String placement) {
|
||||
try {
|
||||
PlacementPolicies.fromName(placement);
|
||||
} catch (IllegalArgumentException e) {
|
||||
throw new IllegalStateException("refusing to start: " + e.getMessage(), e);
|
||||
}
|
||||
}
|
||||
|
||||
static List<String> unknownTopLevelKeys(String yaml) {
|
||||
Map<?, ?> raw;
|
||||
try {
|
||||
|
||||
@@ -577,13 +577,25 @@ public final class BridgeMcp {
|
||||
return text("acknowledged " + msgId);
|
||||
}
|
||||
|
||||
/** {@code bridge_status}: the live lifecycle status of a worker session. */
|
||||
/**
|
||||
* {@code bridge_status}: the live lifecycle status of a worker session, plus — when the worker
|
||||
* is paused mid-turn in an async {@code bridge_ask} (CB-582) — the open question and how to
|
||||
* answer it, so a lead on its normal poll cadence does not need the ticket to notice.
|
||||
*/
|
||||
static McpSchema.CallToolResult status(MessageService messages, String sessionId) {
|
||||
if (isBlank(sessionId)) {
|
||||
return error("sessionId is required");
|
||||
}
|
||||
try {
|
||||
return text(messages.status(sessionId).name().toLowerCase());
|
||||
String base = messages.status(sessionId).name().toLowerCase();
|
||||
MessageService.PendingAsk ask = messages.pendingAsk(sessionId);
|
||||
if (ask == null) {
|
||||
return text(base);
|
||||
}
|
||||
return text(base + "\n\n[question — worker is waiting for your answer]\n" + ask.question()
|
||||
+ "\n\nAnswer it by calling bridge_send again with turnId=\"" + ask.turnId()
|
||||
+ "\" and content set to your answer; the worker resumes the same turn."
|
||||
+ " (ticket " + ask.ticket() + ")");
|
||||
} catch (HerdrException e) {
|
||||
return error("herdr error for session " + sessionId + ": " + e.getMessage());
|
||||
}
|
||||
|
||||
@@ -164,18 +164,30 @@ public final class MessageService {
|
||||
|
||||
/** An in-flight or finished async delegation, keyed by its ticket. */
|
||||
private static final class Task {
|
||||
private final String ticket;
|
||||
private final String target;
|
||||
private final CompletableFuture<Reply> future = new CompletableFuture<>();
|
||||
private final long createdNanos;
|
||||
private volatile Reply question;
|
||||
private volatile String turnId;
|
||||
|
||||
private Task(String target, long createdNanos) {
|
||||
private Task(String ticket, String target, long createdNanos) {
|
||||
this.ticket = ticket;
|
||||
this.target = target;
|
||||
this.createdNanos = createdNanos;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* A worker session's currently-open {@code bridge_ask} question, surfaced so {@code bridge_status}
|
||||
* can show it without the caller needing the ticket first (CB-582). Only covers async
|
||||
* (fire-and-poll) delegations, which track the question on their {@link Task}; a blocking
|
||||
* ({@code wait:true}) send already hands the question straight back to its own caller, so there is
|
||||
* nothing hidden left for {@code bridge_status} to surface in that case.
|
||||
*/
|
||||
public record PendingAsk(String ticket, String question, String turnId) {
|
||||
}
|
||||
|
||||
private final AgentControl agents;
|
||||
private final Injector injector;
|
||||
private final Rendezvous rendezvous;
|
||||
@@ -200,9 +212,11 @@ public final class MessageService {
|
||||
* Create with an explicit {@link ReplyInbox} and optional {@link ReplyPushLoop}.
|
||||
*
|
||||
* @param pushLoop nullable — when non-null, the push loop is notified on the no-waiter reply
|
||||
* branch ({@link #reply}) so it can nudge the primary to drain the inbox, and
|
||||
* branch ({@link #reply}) so it can nudge the primary to drain the inbox,
|
||||
* (CB-588) whenever an async ticket started by {@link #sendAsync} reaches a
|
||||
* terminal phase, and whenever {@link #poll} hands a terminal ticket to its caller
|
||||
* terminal phase, whenever {@link #poll} hands a terminal ticket to its caller,
|
||||
* and (CB-582) whenever an async ticket's worker pauses mid-turn in
|
||||
* {@code bridge_ask} or that pause ends (answered or lapsed)
|
||||
*/
|
||||
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous,
|
||||
ReplyInbox inbox, ReplyPushLoop pushLoop) {
|
||||
@@ -495,6 +509,15 @@ public final class MessageService {
|
||||
rendezvous.closeAsk(ticket.turnId());
|
||||
return new AskResult(AskOutcome.NO_WAITER, null); // no primary is blocked on this worker
|
||||
}
|
||||
// CB-582: the question just became visible via bridge_poll (Phase.ASKING) for an async
|
||||
// (wait:false) delegation — nudge the lead's own pane the same way a terminal ticket does
|
||||
// (CB-588), since the lead's normal poll cadence is minutes away and the reverse-rendezvous
|
||||
// window (~55s, see BridgeMcp/BridgedApp) is far shorter. A blocking (wait:true) send has
|
||||
// no Task and gets the question directly in its own reply, so task == null there — nothing
|
||||
// to nudge.
|
||||
if (task != null && pushLoop != null) {
|
||||
pushLoop.onQuestionOpened(task.ticket, workerSession, ticket.turnId(), question);
|
||||
}
|
||||
}
|
||||
try {
|
||||
String answer = ticket.answer().get(timeoutMillis, TimeUnit.MILLISECONDS);
|
||||
@@ -513,6 +536,14 @@ public final class MessageService {
|
||||
// Only the fresh owner tears down the shared turn; a duplicate must leave it open.
|
||||
if (ticket.fresh()) {
|
||||
rendezvous.closeAsk(ticket.turnId());
|
||||
// CB-582: tear the push loop's copy down at the same point, not only on the three
|
||||
// paths that call clearAsyncQuestion. The answer future can complete exceptionally
|
||||
// (ExecutionException) or the thread be interrupted, and both leave this method by
|
||||
// throwing — the question would stay pending forever, keep being named in nudges
|
||||
// until its own cap, and never be removed from the map. Already-closed is a no-op.
|
||||
if (pushLoop != null) {
|
||||
pushLoop.questionClosed(ticket.turnId());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -589,7 +620,7 @@ public final class MessageService {
|
||||
*/
|
||||
public String sendAsync(String target, String content, Runnable onAccepted) {
|
||||
String ticket = "task-" + ticketSeq.incrementAndGet();
|
||||
Task task = new Task(target, nowNanos.getAsLong());
|
||||
Task task = new Task(ticket, target, nowNanos.getAsLong());
|
||||
tasks.put(ticket, task);
|
||||
if (pushLoop != null) {
|
||||
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
|
||||
@@ -715,6 +746,12 @@ public final class MessageService {
|
||||
|
||||
/** Clear an answered or lapsed question, but only when it matches the ticket's current turn. */
|
||||
private void clearAsyncQuestion(String turnId, boolean forgetTurn) {
|
||||
// CB-582: tell the push loop first — like ticketCollected, a removal for a turnId it never
|
||||
// nudged about (or already dropped) is a harmless no-op, so this is safe to call unconditionally
|
||||
// rather than threading the guard below through it.
|
||||
if (pushLoop != null) {
|
||||
pushLoop.questionClosed(turnId);
|
||||
}
|
||||
Task task = asyncTasksByTurn.get(turnId);
|
||||
if (task != null && turnId.equals(task.turnId)) {
|
||||
task.question = null;
|
||||
@@ -746,6 +783,23 @@ public final class MessageService {
|
||||
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
|
||||
}
|
||||
|
||||
/**
|
||||
* The question {@code workerSession} is currently paused on via {@code bridge_ask}, if any
|
||||
* (CB-582) — {@code bridge_status} uses this to show a pending question without the caller
|
||||
* needing the ticket. {@code null} when the session has no open async question (including a
|
||||
* session mid a <em>blocking</em> {@code bridge_ask}, which has no {@link Task} to look up — see
|
||||
* {@link PendingAsk}).
|
||||
*/
|
||||
public PendingAsk pendingAsk(String workerSession) {
|
||||
for (Task task : tasks.values()) {
|
||||
Reply q = task.question;
|
||||
if (q != null && workerSession.equals(task.target)) {
|
||||
return new PendingAsk(task.ticket, q.text(), q.turnId());
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
/** Release the async executor. */
|
||||
public void close() {
|
||||
asyncExecutor.shutdown();
|
||||
|
||||
@@ -19,30 +19,33 @@ import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
* A status-gated push loop that nudges a lead's own herdr pane when it has uncollected work
|
||||
* waiting: a worker reply queued with no live {@code bridge_send} to resolve it (CB-307), or an
|
||||
* waiting: a worker reply queued with no live {@code bridge_send} to resolve it (CB-307), an
|
||||
* async delegation ticket ({@code bridge_send(wait:false)}) that reached a terminal phase
|
||||
* (CB-588).
|
||||
* (CB-588), or an async ticket's worker pausing mid-turn in {@code bridge_ask} to await an answer
|
||||
* (CB-582).
|
||||
*
|
||||
* <p><strong>CB-590: one schedule per lead.</strong> Both kinds of work are triggered through
|
||||
* their own entry point — {@link #onReplyQueued(String)} and
|
||||
* {@link #onTicketTerminal(String, String, boolean)} — but both resolve the lead that should be
|
||||
* nudged and coalesce onto a single per-lead reminder schedule, tracked in {@link #activeLeads}.
|
||||
* Earlier this was two independent schedules (one keyed by worker target for replies, one keyed
|
||||
* by lead for tickets) that could both decide to inject into the same pane in the same window —
|
||||
* a race, not routine behaviour, but the expensive kind: it interrupts the lead's live turn
|
||||
* twice. Collapsing to one schedule per lead makes that structurally impossible: at most one
|
||||
* scheduled tick chain is ever live for a given lead (guarded by {@link #activeLeads}'
|
||||
* <p><strong>CB-590: one schedule per lead.</strong> All three kinds of work are triggered
|
||||
* through their own entry point — {@link #onReplyQueued(String)},
|
||||
* {@link #onTicketTerminal(String, String, boolean)}, and
|
||||
* {@link #onQuestionOpened(String, String, String, String)} — but each resolves the lead that
|
||||
* should be nudged and coalesces onto a single per-lead reminder schedule, tracked in
|
||||
* {@link #activeLeads}. Earlier this was two independent schedules (one keyed by worker target
|
||||
* for replies, one keyed by lead for tickets) that could both decide to inject into the same pane
|
||||
* in the same window — a race, not routine behaviour, but the expensive kind: it interrupts the
|
||||
* lead's live turn twice. Collapsing to one schedule per lead makes that structurally impossible:
|
||||
* at most one scheduled tick chain is ever live for a given lead (guarded by {@link #activeLeads}'
|
||||
* compare-and-set), so at most one {@code agents.send} to that lead's pane is ever in flight.
|
||||
*
|
||||
* <p>Each tick examines <em>everything</em> pending for that lead — reply targets whose inbox
|
||||
* still holds an unacked message ({@link #pendingReplies}) and tickets not yet collected
|
||||
* ({@link #pendingTickets}) — and sends at most one combined nudge per tick
|
||||
* ({@link #injectNudge(String, int, int)}). Work that arrives while the lead is busy is never
|
||||
* lost: it is re-read fresh on every tick until the lead is injectable or its own reminder cap
|
||||
* ({@link #maxReminders}) is reached — reply and ticket work each spend from their own budget, so
|
||||
* one source exhausting its cap does not stop nudges about the other (post-CB-590 regression fix;
|
||||
* see {@link #decide}) — whichever the durable inbox / pending-ticket set doesn't already answer
|
||||
* via {@code STOP}.
|
||||
* still holds an unacked message ({@link #pendingReplies}), tickets not yet collected
|
||||
* ({@link #pendingTickets}), and open questions not yet answered or lapsed
|
||||
* ({@link #pendingQuestions}) — and sends at most one combined nudge per tick
|
||||
* ({@link #injectNudge(String, int, int, int)}). Work that arrives while the lead is busy is
|
||||
* never lost: it is re-read fresh on every tick until the lead is injectable or its own reminder
|
||||
* cap ({@link #maxReminders}) is reached — each source spends from its own budget, so one source
|
||||
* exhausting its cap does not stop nudges about the others (post-CB-590 regression fix; see
|
||||
* {@link #decide}) — whichever the durable inbox / pending set doesn't already answer via
|
||||
* {@code STOP}.
|
||||
*/
|
||||
public final class ReplyPushLoop {
|
||||
|
||||
@@ -57,6 +60,14 @@ public final class ReplyPushLoop {
|
||||
/** Coalesced form, several uncollected tickets for the same lead. */
|
||||
static final String TICKETS_NUDGE_FORMAT =
|
||||
"%d tickets finished%s — run bridge_poll(ticket=...) for each to collect them: %s";
|
||||
/** Singular form, one worker paused mid-turn in bridge_ask (CB-582) — names the answer call directly. */
|
||||
static final String QUESTION_NUDGE_FORMAT =
|
||||
"Worker %s asked a question (ticket %s) — answer it with bridge_send(turnId=\"%s\", "
|
||||
+ "content=...) to resume its turn:\n%s";
|
||||
/** Coalesced form, several open questions for the same lead. */
|
||||
static final String QUESTIONS_NUDGE_FORMAT =
|
||||
"%d workers are paused on a question — run bridge_poll(ticket=...) for each, then answer "
|
||||
+ "with bridge_send(turnId=..., content=...): %s";
|
||||
|
||||
private final PrimaryRegistry primaryRegistry;
|
||||
private final AgentControl agents;
|
||||
@@ -66,11 +77,23 @@ public final class ReplyPushLoop {
|
||||
private final long backoffMs;
|
||||
private final Metrics metrics; // CB-512: nullable — no registry in unit tests
|
||||
|
||||
/** Worker targets with a reply queued, and the lead to nudge about it, keyed by target. */
|
||||
private final ConcurrentHashMap<String, String> pendingReplies = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* Worker targets with a reply queued, keyed by target. Each entry carries its own nudge
|
||||
* count (CB-598) rather than sharing one counter per lead per source: a target's count only
|
||||
* ever reflects nudges that actually named that target, so a target that joins while the
|
||||
* schedule is already deep into another target's reminders still reads as fresh.
|
||||
*/
|
||||
private final ConcurrentHashMap<String, ReplyEntry> pendingReplies = new ConcurrentHashMap<>();
|
||||
/** Tickets that have gone terminal but not yet been polled, keyed by ticket. */
|
||||
private final ConcurrentHashMap<String, PendingTicket> pendingTickets = new ConcurrentHashMap<>();
|
||||
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets). */
|
||||
/**
|
||||
* Open {@code bridge_ask} questions not yet answered or lapsed, keyed by {@code turnId}
|
||||
* (CB-582). A question's own nudge count is tracked the same per-item way as
|
||||
* {@link #pendingTickets} (CB-598): a fresh question keeps its source eligible regardless of
|
||||
* how depleted an older, still-open question's count is.
|
||||
*/
|
||||
private final ConcurrentHashMap<String, PendingQuestion> pendingQuestions = new ConcurrentHashMap<>();
|
||||
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets and/or questions). */
|
||||
private final ConcurrentHashMap<String, Boolean> activeLeads = new ConcurrentHashMap<>();
|
||||
|
||||
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||
@@ -116,10 +139,10 @@ public final class ReplyPushLoop {
|
||||
Set<String> result = new HashSet<>();
|
||||
for (var entry : pendingReplies.entrySet()) {
|
||||
String target = entry.getKey();
|
||||
String owningLead = entry.getValue();
|
||||
if (!lead.equals(owningLead)) continue;
|
||||
ReplyEntry owning = entry.getValue();
|
||||
if (!lead.equals(owning.lead())) continue;
|
||||
if (inbox.peek(target).isEmpty()) {
|
||||
pendingReplies.remove(target, owningLead);
|
||||
pendingReplies.remove(target, owning);
|
||||
continue;
|
||||
}
|
||||
result.add(target);
|
||||
@@ -127,8 +150,15 @@ public final class ReplyPushLoop {
|
||||
return result;
|
||||
}
|
||||
|
||||
/** A ticket awaiting collection: which lead to nudge, and whether it ended in failure. */
|
||||
private record PendingTicket(String ticket, String lead, boolean failed) {
|
||||
/** A pending reply target: which lead to nudge, and how many nudges have named it so far. */
|
||||
private record ReplyEntry(String lead, int nudgeCount) {
|
||||
}
|
||||
|
||||
/**
|
||||
* A ticket awaiting collection: which lead to nudge, whether it ended in failure, and how
|
||||
* many nudges have named it so far (CB-598 — tracked per ticket, not per lead per source).
|
||||
*/
|
||||
private record PendingTicket(String ticket, String lead, boolean failed, int nudgeCount) {
|
||||
}
|
||||
|
||||
/** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */
|
||||
@@ -142,6 +172,70 @@ public final class ReplyPushLoop {
|
||||
.collect(Collectors.toUnmodifiableSet());
|
||||
}
|
||||
|
||||
/**
|
||||
* An open question awaiting the lead's answer: which ticket it belongs to, which worker asked,
|
||||
* which lead to nudge, the question text, and how many nudges have named it so far (CB-598 —
|
||||
* tracked per question, not per lead per source).
|
||||
*/
|
||||
private record PendingQuestion(String turnId, String ticket, String target, String lead,
|
||||
String question, int nudgeCount) {
|
||||
}
|
||||
|
||||
/** Questions still open for {@code lead}, snapshotted fresh for one tick. */
|
||||
private List<PendingQuestion> pendingQuestionsFor(String lead) {
|
||||
return pendingQuestions.values().stream().filter(q -> lead.equals(q.lead())).toList();
|
||||
}
|
||||
|
||||
/** Question turnIds still open for {@code lead} — a plain snapshot for race comparison. */
|
||||
private Set<String> pendingQuestionTurnIdsFor(String lead) {
|
||||
return pendingQuestionsFor(lead).stream().map(PendingQuestion::turnId)
|
||||
.collect(Collectors.toUnmodifiableSet());
|
||||
}
|
||||
|
||||
/**
|
||||
* The reply-source reminder count {@link #decide} should see for {@code lead} on this tick:
|
||||
* the <em>minimum</em> nudge count among the reply targets currently pending for it (CB-598).
|
||||
*
|
||||
* <p>Before this, the count passed to {@code decide} was a single counter carried forward
|
||||
* across scheduled ticks ({@code scheduleNext(lead, count + 1, ...)}), incremented whenever
|
||||
* the source had <em>any</em> pending work — not tied to which target that work was. A target
|
||||
* that joined while an older target's count was already near the cap inherited that count on
|
||||
* its very next tick, even though no nudge had ever named it. Taking the minimum over what is
|
||||
* actually pending now means a fresh target (count 0) keeps the source eligible regardless of
|
||||
* how many times an older, still-undrained target has already been nudged; that older target
|
||||
* keeps riding along in the combined nudge text without spending any more of its own budget
|
||||
* (see {@link #bumpNudgeCounts}). Returns 0 when nothing is pending — {@link #decide} never
|
||||
* consults the count in that case, since {@code hasReplyWork} is false.
|
||||
*/
|
||||
private int minReplyNudgeCountFor(String lead) {
|
||||
int min = Integer.MAX_VALUE;
|
||||
for (String target : pendingReplyTargetsFor(lead)) {
|
||||
ReplyEntry entry = pendingReplies.get(target);
|
||||
if (entry != null) {
|
||||
min = Math.min(min, entry.nudgeCount());
|
||||
}
|
||||
}
|
||||
return min == Integer.MAX_VALUE ? 0 : min;
|
||||
}
|
||||
|
||||
/** As {@link #minReplyNudgeCountFor}, for the ticket source. */
|
||||
private int minTicketNudgeCountFor(String lead) {
|
||||
int min = Integer.MAX_VALUE;
|
||||
for (PendingTicket ticket : pendingTicketsFor(lead)) {
|
||||
min = Math.min(min, ticket.nudgeCount());
|
||||
}
|
||||
return min == Integer.MAX_VALUE ? 0 : min;
|
||||
}
|
||||
|
||||
/** As {@link #minReplyNudgeCountFor}, for the question source (CB-582). */
|
||||
private int minQuestionNudgeCountFor(String lead) {
|
||||
int min = Integer.MAX_VALUE;
|
||||
for (PendingQuestion q : pendingQuestionsFor(lead)) {
|
||||
min = Math.min(min, q.nudgeCount());
|
||||
}
|
||||
return min == Integer.MAX_VALUE ? 0 : min;
|
||||
}
|
||||
|
||||
/**
|
||||
* Pure decision function: examine everything pending for {@code lead} — reply targets and
|
||||
* tickets alike — and return what the loop should do.
|
||||
@@ -155,21 +249,40 @@ public final class ReplyPushLoop {
|
||||
* {@link Action#INJECT}. Only when neither source has eligible work does the loop
|
||||
* {@link Action#STOP}.
|
||||
*
|
||||
* <p><strong>CB-598: the counts are per-item, not per-tick.</strong> {@link #tick} no longer
|
||||
* carries these counts forward across scheduled calls — it recomputes them fresh every tick via
|
||||
* {@link #minReplyNudgeCountFor} / {@link #minTicketNudgeCountFor}, so this function itself did
|
||||
* not need to change; only what its caller feeds it did.
|
||||
*
|
||||
* @param lead the lead terminal to nudge
|
||||
* @param replyReminderCount how many nudges have covered pending reply work for this lead
|
||||
* @param ticketReminderCount how many nudges have covered pending ticket work for this lead
|
||||
* @param replyReminderCount the lowest nudge count among reply targets pending for this lead
|
||||
* @param ticketReminderCount the lowest nudge count among tickets pending for this lead
|
||||
* @return the action the caller should take
|
||||
*/
|
||||
Action decide(String lead, int replyReminderCount, int ticketReminderCount) {
|
||||
return decide(lead, replyReminderCount, ticketReminderCount, minQuestionNudgeCountFor(lead));
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #decide(String, int, int)}, with the question source (CB-582) folded in on the
|
||||
* same footing as replies and tickets: its own eligibility (has open questions AND under its
|
||||
* own {@link #maxReminders} budget) is enough on its own to {@link Action#INJECT}, exactly like
|
||||
* the other two.
|
||||
*
|
||||
* @param questionReminderCount the lowest nudge count among questions open for this lead
|
||||
*/
|
||||
Action decide(String lead, int replyReminderCount, int ticketReminderCount, int questionReminderCount) {
|
||||
boolean hasReplyWork = !pendingReplyTargetsFor(lead).isEmpty();
|
||||
boolean hasTicketWork = !pendingTicketIdsFor(lead).isEmpty();
|
||||
if (!hasReplyWork && !hasTicketWork) {
|
||||
boolean hasQuestionWork = !pendingQuestionTurnIdsFor(lead).isEmpty();
|
||||
if (!hasReplyWork && !hasTicketWork && !hasQuestionWork) {
|
||||
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
|
||||
return Action.STOP;
|
||||
}
|
||||
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
|
||||
boolean ticketEligible = hasTicketWork && ticketReminderCount < maxReminders;
|
||||
if (!replyEligible && !ticketEligible) {
|
||||
boolean questionEligible = hasQuestionWork && questionReminderCount < maxReminders;
|
||||
if (!replyEligible && !ticketEligible && !questionEligible) {
|
||||
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
|
||||
maxReminders, lead);
|
||||
countNudge("exhausted");
|
||||
@@ -204,7 +317,8 @@ public final class ReplyPushLoop {
|
||||
log.debug("push: no lead is known to be waiting on {}, skipping reminder", target);
|
||||
return;
|
||||
}
|
||||
pendingReplies.put(target, lead.get());
|
||||
pendingReplies.compute(target, (t, existing) ->
|
||||
new ReplyEntry(lead.get(), existing == null ? 0 : existing.nudgeCount()));
|
||||
startOrCoalesce(lead.get());
|
||||
}
|
||||
|
||||
@@ -232,7 +346,8 @@ public final class ReplyPushLoop {
|
||||
ticket, target);
|
||||
return;
|
||||
}
|
||||
pendingTickets.put(ticket, new PendingTicket(ticket, lead.get(), failed));
|
||||
pendingTickets.compute(ticket, (id, existing) ->
|
||||
new PendingTicket(ticket, lead.get(), failed, existing == null ? 0 : existing.nudgeCount()));
|
||||
startOrCoalesce(lead.get());
|
||||
}
|
||||
|
||||
@@ -246,6 +361,39 @@ public final class ReplyPushLoop {
|
||||
pendingTickets.remove(ticket);
|
||||
}
|
||||
|
||||
/**
|
||||
* Called when an async ticket's worker pauses mid-turn in {@code bridge_ask} (CB-582): the
|
||||
* question is now visible via {@code bridge_poll} (Phase.ASKING), but the reverse-rendezvous
|
||||
* window it opened with (~55s default, see {@code BridgeMcp}/{@code BridgedApp}) is far shorter
|
||||
* than a lead's normal minutes-long poll cadence — exactly the gap this closes. Resolves the
|
||||
* delegating lead the same way {@link #onTicketTerminal} does and coalesces onto the same
|
||||
* per-lead schedule (CB-590).
|
||||
*
|
||||
* @param ticket the async ticket the question belongs to (for {@code bridge_poll})
|
||||
* @param target the worker session that asked
|
||||
* @param turnId correlation id the lead answers with ({@code bridge_send turnId=...})
|
||||
* @param question the question text
|
||||
*/
|
||||
public void onQuestionOpened(String ticket, String target, String turnId, String question) {
|
||||
var lead = primaryRegistry.nudgeTargetFor(target);
|
||||
if (lead.isEmpty()) {
|
||||
log.debug("push: no lead is known to be waiting on {}'s question (turnId {}), skipping nudge",
|
||||
target, turnId);
|
||||
return;
|
||||
}
|
||||
pendingQuestions.put(turnId, new PendingQuestion(turnId, ticket, target, lead.get(), question, 0));
|
||||
startOrCoalesce(lead.get());
|
||||
}
|
||||
|
||||
/**
|
||||
* Called when a worker's {@code bridge_ask} resolves — answered or lapsed unanswered — so a
|
||||
* scheduled tick never nudges about a question the lead already handled. A {@code turnId} that
|
||||
* was never pending (never nudged, or already closed) is a no-op.
|
||||
*/
|
||||
public void questionClosed(String turnId) {
|
||||
pendingQuestions.remove(turnId);
|
||||
}
|
||||
|
||||
// --- the schedule ----------------------------------------------------------------------------
|
||||
|
||||
/** Start a reminder schedule for {@code lead}, or join the one already running. */
|
||||
@@ -255,29 +403,40 @@ public final class ReplyPushLoop {
|
||||
return;
|
||||
}
|
||||
log.debug("push: starting reminder loop for lead {}", lead);
|
||||
scheduleNext(lead, 0, 0);
|
||||
scheduleNext(lead);
|
||||
}
|
||||
|
||||
/** Execute one loop tick — called on the scheduler thread. */
|
||||
private void tick(String lead, int replyReminderCount, int ticketReminderCount) {
|
||||
/**
|
||||
* Execute one loop tick — called on the scheduler thread (or directly by a test; package-private
|
||||
* for the same reason as {@link #stopOrRestart}).
|
||||
*
|
||||
* <p><strong>CB-598.</strong> The reminder counts fed into {@link #decide} are recomputed fresh
|
||||
* every tick from what is actually pending right now ({@link #minReplyNudgeCountFor} /
|
||||
* {@link #minTicketNudgeCountFor}), rather than carried forward as running counters across
|
||||
* scheduled calls. A counter carried forward has no memory of which item it was counting for:
|
||||
* a target or ticket that joined mid-backoff — after the previous tick fired but before this one
|
||||
* did — is already sitting in {@code repliesBefore} / {@code ticketsBefore} below by the time this
|
||||
* tick takes its snapshot, indistinguishable at that point from backlog the cap is meant to
|
||||
* silence. Recomputing from the per-item counts fixes that: a newly-joined item's own count is
|
||||
* still 0, so it keeps its source eligible regardless of how depleted an older, still-undrained
|
||||
* item's count is.
|
||||
*/
|
||||
void tick(String lead) {
|
||||
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
|
||||
Set<String> ticketsBefore = pendingTicketIdsFor(lead);
|
||||
var action = decide(lead, replyReminderCount, ticketReminderCount);
|
||||
Set<String> questionsBefore = pendingQuestionTurnIdsFor(lead);
|
||||
int replyReminderCount = minReplyNudgeCountFor(lead);
|
||||
int ticketReminderCount = minTicketNudgeCountFor(lead);
|
||||
int questionReminderCount = minQuestionNudgeCountFor(lead);
|
||||
var action = decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount);
|
||||
switch (action) {
|
||||
case INJECT -> {
|
||||
injectNudge(lead, replyReminderCount, ticketReminderCount);
|
||||
// Only the source(s) actually eligible this tick spend a unit of their own budget —
|
||||
// an exhausted source riding along in the combined message (still pending, still
|
||||
// named) does not get charged again; its count stays put until it drains.
|
||||
boolean replyEligible = !repliesBefore.isEmpty() && replyReminderCount < maxReminders;
|
||||
boolean ticketEligible = !ticketsBefore.isEmpty() && ticketReminderCount < maxReminders;
|
||||
scheduleNext(lead,
|
||||
replyEligible ? replyReminderCount + 1 : replyReminderCount,
|
||||
ticketEligible ? ticketReminderCount + 1 : ticketReminderCount);
|
||||
injectNudge(lead, replyReminderCount, ticketReminderCount, questionReminderCount);
|
||||
scheduleNext(lead);
|
||||
}
|
||||
// Re-check after the configured backoff; the lead may become injectable soon.
|
||||
case WAIT_BUSY -> scheduleNext(lead, replyReminderCount, ticketReminderCount);
|
||||
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore);
|
||||
case WAIT_BUSY -> scheduleNext(lead);
|
||||
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -312,50 +471,92 @@ public final class ReplyPushLoop {
|
||||
* side always wins; neither can miss the other, so this never loops on its own account.
|
||||
*/
|
||||
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore) {
|
||||
stopOrRestart(lead, repliesBefore, ticketsBefore, pendingQuestionTurnIdsFor(lead));
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #stopOrRestart(String, Set, Set)}, with the question source's (CB-582) own "before"
|
||||
* snapshot folded into the same race check: a question that raced in during the
|
||||
* decision-to-release window reclaims the schedule slot exactly like a raced-in reply or ticket.
|
||||
*/
|
||||
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore,
|
||||
Set<String> questionsBefore) {
|
||||
activeLeads.remove(lead);
|
||||
boolean racedIn = pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t))
|
||||
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t));
|
||||
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t))
|
||||
|| pendingQuestionTurnIdsFor(lead).stream().anyMatch(t -> !questionsBefore.contains(t));
|
||||
if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
|
||||
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
|
||||
scheduleNext(lead, 0, 0);
|
||||
scheduleNext(lead);
|
||||
return;
|
||||
}
|
||||
log.debug("push: reminder loop ended for lead {}", lead);
|
||||
}
|
||||
|
||||
/** Send one combined nudge covering everything currently pending for {@code lead}. */
|
||||
private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount) {
|
||||
// Re-read rather than threading it down from decide(): a reply can drain, or a ticket be
|
||||
// collected (or another arrive), between the decision and the injection.
|
||||
private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount,
|
||||
int questionReminderCount) {
|
||||
// Re-read rather than threading it down from decide(): a reply can drain, a ticket be
|
||||
// collected, or a question be answered (or another arrive), between the decision and the
|
||||
// injection.
|
||||
Set<String> replyTargets = pendingReplyTargetsFor(lead);
|
||||
List<PendingTicket> tickets = pendingTicketsFor(lead);
|
||||
if (replyTargets.isEmpty() && tickets.isEmpty()) {
|
||||
List<PendingQuestion> questions = pendingQuestionsFor(lead);
|
||||
if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty()) {
|
||||
log.debug("push: pending work for lead {} drained before the nudge could be sent", lead);
|
||||
return;
|
||||
}
|
||||
String nudge = formatNudge(replyTargets, tickets);
|
||||
String nudge = formatNudge(replyTargets, tickets, questions);
|
||||
try {
|
||||
agents.send(lead, nudge);
|
||||
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}; {} reply target(s), {} ticket(s))",
|
||||
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; "
|
||||
+ "{} reply target(s), {} ticket(s), {} question(s))",
|
||||
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
|
||||
replyTargets.size(), tickets.size());
|
||||
questionReminderCount + 1, maxReminders,
|
||||
replyTargets.size(), tickets.size(), questions.size());
|
||||
countNudge("delivered");
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}): {}",
|
||||
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, e.toString());
|
||||
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}, question {}/{}): {}",
|
||||
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
|
||||
questionReminderCount + 1, maxReminders, e.toString());
|
||||
}
|
||||
// Bump every item actually named in this nudge, not just whatever the shared source-level
|
||||
// eligibility used to gate (CB-598) — each item's own count is what the next tick's
|
||||
// minReplyNudgeCountFor / minTicketNudgeCountFor / minQuestionNudgeCountFor will read. An
|
||||
// item already at or over the cap keeps riding along in the text (still pending, still
|
||||
// named) but its extra bumps here are inert: decide() already treats it as ineligible once
|
||||
// its count reaches maxReminders.
|
||||
bumpNudgeCounts(replyTargets, tickets, questions);
|
||||
}
|
||||
|
||||
/** Record that every one of these items was just named in a sent (or attempted) nudge. */
|
||||
private void bumpNudgeCounts(Set<String> replyTargets, List<PendingTicket> tickets,
|
||||
List<PendingQuestion> questions) {
|
||||
for (String target : replyTargets) {
|
||||
pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1));
|
||||
}
|
||||
for (PendingTicket ticket : tickets) {
|
||||
pendingTickets.computeIfPresent(ticket.ticket(),
|
||||
(id, e) -> new PendingTicket(e.ticket(), e.lead(), e.failed(), e.nudgeCount() + 1));
|
||||
}
|
||||
for (PendingQuestion question : questions) {
|
||||
pendingQuestions.computeIfPresent(question.turnId(), (id, e) ->
|
||||
new PendingQuestion(e.turnId(), e.ticket(), e.target(), e.lead(), e.question(),
|
||||
e.nudgeCount() + 1));
|
||||
}
|
||||
}
|
||||
|
||||
/** Schedule the next tick on the scheduler thread pool. */
|
||||
private void scheduleNext(String lead, int nextReplyReminderCount, int nextTicketReminderCount) {
|
||||
scheduler.schedule(() -> tick(lead, nextReplyReminderCount, nextTicketReminderCount),
|
||||
private void scheduleNext(String lead) {
|
||||
scheduler.schedule(() -> tick(lead),
|
||||
backoffMs, TimeUnit.MILLISECONDS);
|
||||
}
|
||||
|
||||
// --- nudge formatting ------------------------------------------------------------------------
|
||||
|
||||
/** Render everything pending for one lead as a single nudge line. */
|
||||
private static String formatNudge(Set<String> replyTargets, List<PendingTicket> tickets) {
|
||||
private static String formatNudge(Set<String> replyTargets, List<PendingTicket> tickets,
|
||||
List<PendingQuestion> questions) {
|
||||
List<String> parts = new ArrayList<>();
|
||||
if (!replyTargets.isEmpty()) {
|
||||
parts.add(formatRepliesNudge(replyTargets));
|
||||
@@ -363,6 +564,9 @@ public final class ReplyPushLoop {
|
||||
if (!tickets.isEmpty()) {
|
||||
parts.add(formatTicketsNudge(tickets));
|
||||
}
|
||||
if (!questions.isEmpty()) {
|
||||
parts.add(formatQuestionsNudge(questions));
|
||||
}
|
||||
return String.join(" | ", parts);
|
||||
}
|
||||
|
||||
@@ -390,6 +594,18 @@ public final class ReplyPushLoop {
|
||||
return TICKETS_NUDGE_FORMAT.formatted(pending.size(), failedNote, ids);
|
||||
}
|
||||
|
||||
/** Render one or several open questions (CB-582). */
|
||||
private static String formatQuestionsNudge(List<PendingQuestion> pending) {
|
||||
if (pending.size() == 1) {
|
||||
PendingQuestion q = pending.get(0);
|
||||
return QUESTION_NUDGE_FORMAT.formatted(q.target(), q.ticket(), q.turnId(), q.question());
|
||||
}
|
||||
String ids = pending.stream()
|
||||
.map(q -> q.ticket() + " (turnId=" + q.turnId() + ")")
|
||||
.collect(Collectors.joining(", "));
|
||||
return QUESTIONS_NUDGE_FORMAT.formatted(pending.size(), ids);
|
||||
}
|
||||
|
||||
// --- lifecycle -----------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
@@ -397,8 +613,8 @@ public final class ReplyPushLoop {
|
||||
* uses this to stand aside: while the push loop is actively nudging a lead, a concurrent
|
||||
* heartbeat injection would start a second competing turn in the same pane — racing loops
|
||||
* multiply turns and context burn. "Active" means a schedule exists in {@link #activeLeads},
|
||||
* which now covers both reply-queued (CB-307) and ticket-terminal (CB-588) work (CB-590) —
|
||||
* bounded by what has been triggered, not by any persistent state.
|
||||
* which now covers reply-queued (CB-307), ticket-terminal (CB-588), and question-open (CB-582)
|
||||
* work (CB-590) — bounded by what has been triggered, not by any persistent state.
|
||||
*/
|
||||
public boolean isActive() {
|
||||
return !activeLeads.isEmpty();
|
||||
@@ -410,6 +626,7 @@ public final class ReplyPushLoop {
|
||||
activeLeads.clear();
|
||||
pendingReplies.clear();
|
||||
pendingTickets.clear();
|
||||
pendingQuestions.clear();
|
||||
}
|
||||
|
||||
/** @see #stop() */
|
||||
|
||||
@@ -509,10 +509,20 @@ public final class BridgedApp {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
ctx.status(200).json(Map.of(
|
||||
"sessionId", id,
|
||||
"status", messages.status(id).name().toLowerCase(),
|
||||
"ready", presence.isPresent(id)));
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
body.put("sessionId", id);
|
||||
body.put("status", messages.status(id).name().toLowerCase());
|
||||
body.put("ready", presence.isPresent(id));
|
||||
// CB-582: a worker paused mid-turn in an async bridge_ask is otherwise invisible to a
|
||||
// status poll — surface the open question and how to answer it, same as bridge_poll's
|
||||
// Phase.ASKING view.
|
||||
MessageService.PendingAsk ask = messages.pendingAsk(id);
|
||||
if (ask != null) {
|
||||
body.put("question", ask.question());
|
||||
body.put("turnId", ask.turnId());
|
||||
body.put("ticket", ask.ticket());
|
||||
}
|
||||
ctx.status(200).json(body);
|
||||
} catch (HerdrException e) {
|
||||
herdrError(ctx, e);
|
||||
}
|
||||
@@ -538,6 +548,11 @@ public final class BridgedApp {
|
||||
if (v.detail() != null) {
|
||||
body.put("detail", v.detail());
|
||||
}
|
||||
// CB-582: Phase.ASKING carries the question in v.reply() (handled above) and its answer-
|
||||
// correlation id here — a REST caller polling this ticket otherwise has no way to answer it.
|
||||
if (v.turnId() != null) {
|
||||
body.put("turnId", v.turnId());
|
||||
}
|
||||
ctx.status(200).json(body);
|
||||
}
|
||||
|
||||
|
||||
@@ -293,8 +293,10 @@ public final class SessionManager implements TurnListener {
|
||||
// blocked caller fails fast with a real reason instead of sitting on a rendezvous
|
||||
// nothing will ever resolve. CB-578 stage C: carry the worktree/branch/snapshot ref
|
||||
// too, so a failed ticket's detail can point a lead at the same tree to re-dispatch.
|
||||
// CB-584 (issue #65 criterion 5): carry agentSessionId alongside them, so a lead can
|
||||
// also resume the member's conversation, not just re-dispatch onto its files.
|
||||
notifyReleased(new ReleaseDetail(removed.terminalId(), removed.worktree(),
|
||||
removed.branch(), snapshotRef));
|
||||
removed.branch(), snapshotRef, removed.agentSessionId()));
|
||||
}
|
||||
}
|
||||
// CB-581: the pane must always stop, even if the dirty check above threw. A session removed
|
||||
@@ -348,7 +350,8 @@ public final class SessionManager implements TurnListener {
|
||||
* are {@code null} for a shared-tree session; {@code snapshotRef} is {@code null} unless this
|
||||
* release snapshotted a dirty worktree into {@code refs/wip/<branch>}.
|
||||
*/
|
||||
public record ReleaseDetail(String terminalId, String worktreePath, String branch, String snapshotRef) {
|
||||
public record ReleaseDetail(String terminalId, String worktreePath, String branch, String snapshotRef,
|
||||
String agentSessionId) {
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -1040,6 +1040,44 @@ class BridgedConfigTest {
|
||||
"an opencode worker with no argv defaults to the opencode binary, never claude");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-604: an unrecognized {@code kind:} used to silently fall into the claude-code bucket — not
|
||||
* matching {@code "opencode"} was the only check. With {@code argv:} also unset that meant the
|
||||
* daemon tried to launch a program literally named after the typo.
|
||||
*/
|
||||
@Test
|
||||
void unknownKindIsRefusedAtLoadNamingTheValueAndTheAcceptedSet(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("kind-typo.yaml");
|
||||
Files.writeString(f, """
|
||||
profiles:
|
||||
gemini:
|
||||
kind: opencod
|
||||
model: google/gemini-2.5-pro
|
||||
""");
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("gemini"), "error names the profile: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("opencod"), "error names the bad value: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("claude-code") && e.getMessage().contains("opencode"),
|
||||
"error names the accepted set: " + e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void blankKindStillDefaultsToClaudeCode(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("kind-blank.yaml");
|
||||
Files.writeString(f, """
|
||||
profiles:
|
||||
gx10:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
kind: ""
|
||||
argv: ["ccs", "gx10"]
|
||||
""");
|
||||
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
assertEquals(BridgedConfig.Profile.KIND_CLAUDE_CODE, cfg.profiles().get("gx10").kind(),
|
||||
"a blank kind: is documented to behave exactly like an absent one");
|
||||
}
|
||||
|
||||
@Test
|
||||
void authDefaultsToLoopbackTrustSoExistingConfigsBehaveAsBefore(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-auth-block.yaml");
|
||||
@@ -1369,6 +1407,89 @@ class BridgedConfigTest {
|
||||
assertTrue(e.getMessage().contains("maxLoad"), "error names the key: " + e.getMessage());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-606: an unrecognized {@code auth.mode} used to silently fall back to
|
||||
* {@code loopback-trust} — {@link BridgedConfig.Auth#tokenMode()} only checked equality
|
||||
* against {@code "token"}. On a loopback bind {@link BridgedConfig#validateAuthExposure()}
|
||||
* never runs (it only fires for a non-loopback bind), so the typo was completely invisible:
|
||||
* the daemon started cleanly and authenticated nobody while the operator believed token mode
|
||||
* was active.
|
||||
*/
|
||||
@Test
|
||||
void unknownAuthModeIsRefusedAtLoadNamingTheValueAndTheAcceptedSet(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("auth-mode-typo.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
host: 127.0.0.1
|
||||
port: 8765
|
||||
auth:
|
||||
mode: toekn
|
||||
""");
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("toekn"), "error names the bad value: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("loopback-trust") && e.getMessage().contains("token"),
|
||||
"error names the accepted set: " + e.getMessage());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-606: an unrecognized per-profile {@code placement:} used to silently fall back to legacy
|
||||
* pane placement — {@link BridgedConfig.Profile#tabPlacement()} only checked equality against
|
||||
* {@code "tab"}.
|
||||
*/
|
||||
@Test
|
||||
void unknownProfilePlacementIsRefusedAtLoadNamingTheProfileAndTheAcceptedSet(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("placement-typo.yaml");
|
||||
Files.writeString(f, """
|
||||
profiles:
|
||||
gx10:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
placement: tabb
|
||||
""");
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("gx10"), "error names the profile: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("tabb"), "error names the bad value: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("tab") && e.getMessage().contains("pane"),
|
||||
"error names the accepted set: " + e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void absentProfilePlacementDefaultsToTab(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("placement-absent.yaml");
|
||||
Files.writeString(f, """
|
||||
profiles:
|
||||
gx10:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
""");
|
||||
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
BridgedConfig.Profile w = cfg.profiles().get("gx10");
|
||||
assertTrue(w.tabPlacement(), "an absent placement must keep defaulting to tab");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-606: the top-level {@code placement:} policy name WAS validated, but only lazily, by
|
||||
* {@code PlacementPolicies.fromName} through {@code CompositePeerLauncher}'s per-spawn
|
||||
* {@code Supplier} — so a bad name still started a daemon that looked healthy and failed only
|
||||
* the first time something spawned without naming a profile. This must now fail at load.
|
||||
*/
|
||||
@Test
|
||||
void unknownTopLevelPlacementPolicyIsRefusedAtLoadNotLazilyAtFirstSpawn(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("placement-policy-typo.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
placement: weightd
|
||||
""");
|
||||
|
||||
IllegalStateException e = assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("weightd"), "error names the bad value: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("fixed") && e.getMessage().contains("round-robin")
|
||||
&& e.getMessage().contains("weighted"),
|
||||
"error names the accepted set: " + e.getMessage());
|
||||
}
|
||||
|
||||
@Test
|
||||
void subscriptionFlagBindsAndDefaultsFalse(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("subscription.yaml");
|
||||
|
||||
@@ -7,6 +7,7 @@ import java.util.ArrayList;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
|
||||
/**
|
||||
* Recording fake {@link HerdrClient} for unit/acceptance tests. Returns canned frames
|
||||
@@ -22,7 +23,13 @@ public final class FakeHerdr implements HerdrClient {
|
||||
public static final long WORKER_PID = 4242;
|
||||
|
||||
private final ObjectMapper mapper = new ObjectMapper();
|
||||
public final List<Call> calls = new ArrayList<>();
|
||||
/**
|
||||
* Thread-safe on purpose. Background loops — {@link dev.ltms.bridged.msg.ReplyPushLoop} and the
|
||||
* lead heartbeat — call this fake from their own scheduler threads while a test polls
|
||||
* {@link #called} from the test thread. A plain {@code ArrayList} threw
|
||||
* {@code ConcurrentModificationException} out of {@code called()} when a nudge landed mid-stream.
|
||||
*/
|
||||
public final List<Call> calls = new CopyOnWriteArrayList<>();
|
||||
private boolean healthy = true;
|
||||
private final List<String> extraWorkspaces = new ArrayList<>();
|
||||
private final List<String> extraAgents = new ArrayList<>();
|
||||
|
||||
@@ -773,6 +773,52 @@ class BridgeMcpTest {
|
||||
assertEquals("blocked", textOf(res));
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-582: a lead polling {@code bridge_status} on its normal cadence — not {@code bridge_poll}
|
||||
* — must also see a worker's open async {@code bridge_ask} question, since the reverse-rendezvous
|
||||
* window it opened with is far shorter than that cadence.
|
||||
*/
|
||||
@Test
|
||||
void statusReportsAnOpenQuestionWhenTheWorkerIsMidAsk() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(T), "sendAsync should have opened its rendezvous waiter");
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
||||
|
||||
MessageService.TaskView asking;
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
do {
|
||||
asking = messages.poll(ticket);
|
||||
Thread.sleep(5);
|
||||
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
|
||||
assertEquals(MessageService.Phase.ASKING, asking.phase());
|
||||
|
||||
McpSchema.CallToolResult res = BridgeMcp.status(messages, T);
|
||||
assertNotEquals(Boolean.TRUE, res.isError());
|
||||
String out = textOf(res);
|
||||
assertTrue(out.startsWith("idle"), "the live status must still lead the text: " + out);
|
||||
assertTrue(out.contains("which config file?"), "the question text must be shown: " + out);
|
||||
assertTrue(out.contains("turnId=\"" + asking.turnId() + "\""), "the turnId must be shown: " + out);
|
||||
assertTrue(out.contains("(ticket " + ticket + ")"), "the ticket must be shown: " + out);
|
||||
|
||||
// Clean up the still-open ask so the background thread does not linger past the test.
|
||||
String turnId = asking.turnId();
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(turnId, "config.yaml", 5000));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
answer.get(5, TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
// --- bridge_whoami: the caller's own identity, so an agent never has to guess its role -------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -35,9 +35,12 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
*/
|
||||
class AmqpReplyInboxRecoveryRaceTest {
|
||||
|
||||
/** Large enough that the (unfixed) unsynchronized sweep's iteration is a real, observable window
|
||||
* a concurrently-started publish can land in — not just a best case, single-entry sprint. */
|
||||
private static final int STALE_PUBLISHES = 100_000;
|
||||
/** Large enough that thousands of entries are still unprocessed by the time the very first one
|
||||
* is observed as failed (see {@code sweepIsHoldingTheLock} below) — that gap is what makes the
|
||||
* head start deterministic instead of a coin flip. 100,000 gave the same guarantee but made the
|
||||
* test far more expensive than the guarantee needs; the ordering no longer depends on a timing
|
||||
* window sized to the full backlog; just to the tail of it. */
|
||||
private static final int STALE_PUBLISHES = 2_000;
|
||||
|
||||
@Test
|
||||
@Timeout(30)
|
||||
@@ -58,6 +61,12 @@ class AmqpReplyInboxRecoveryRaceTest {
|
||||
// pendingByMsgId exactly like publishes whose confirm never arrived before a connection drop.
|
||||
// Virtual threads make this many concurrent blocking publish() calls cheap.
|
||||
CountDownLatch staleStarted = new CountDownLatch(STALE_PUBLISHES);
|
||||
// Counted down by the FIRST stale publish thread to observe its own failure. That can only
|
||||
// happen from inside failPendingPublishesOnRecovery() — nothing else in this test ever
|
||||
// completes a stale Pending exceptionally (no nack/return is simulated for any "stale-*"
|
||||
// msgId) — so seeing it fire is direct, observable proof the sweep is inside its loop, not a
|
||||
// timing guess. It replaces the old fixed Thread.sleep(5) head start.
|
||||
CountDownLatch sweepIsHoldingTheLock = new CountDownLatch(1);
|
||||
for (int i = 0; i < STALE_PUBLISHES; i++) {
|
||||
String msgId = "stale-" + i;
|
||||
Thread.ofVirtual().start(() -> {
|
||||
@@ -65,7 +74,7 @@ class AmqpReplyInboxRecoveryRaceTest {
|
||||
try {
|
||||
inbox.publish("worker-stale", msgId, "x");
|
||||
} catch (IllegalStateException expected) {
|
||||
// resolved (failed by the sweep) — that is exactly what this thread is here for
|
||||
sweepIsHoldingTheLock.countDown();
|
||||
}
|
||||
});
|
||||
}
|
||||
@@ -83,13 +92,28 @@ class AmqpReplyInboxRecoveryRaceTest {
|
||||
}
|
||||
}, "recovery-sweep");
|
||||
sweepThread.start();
|
||||
// A short, deliberate head start: with STALE_PUBLISHES this large, the (unfixed) sweep's own
|
||||
// iteration takes several milliseconds, so this guarantees the sweep has already begun —
|
||||
// and, once guarded, is already holding publishChannelLock — before "fresh" attempts to
|
||||
// register. Without this head start, "fresh" sometimes wins the race for the lock and
|
||||
// registers before the sweep even starts, which is the accepted "already in flight when
|
||||
// recovery fires" case (correctly failed either way) rather than the bug under test.
|
||||
Thread.sleep(5);
|
||||
|
||||
// Deterministic head start: block until the sweep has actually failed one of the stale
|
||||
// publishes. failPendingPublishesOnRecovery() (once guarded, as it is on main) holds
|
||||
// publishChannelLock for its ENTIRE loop, not just per entry — so this failure proves the
|
||||
// sweep is, at this instant, still holding that lock. With STALE_PUBLISHES this large, the
|
||||
// remaining ~1,999 entries give an enormous margin between "first failure observed" and "sweep
|
||||
// releases the lock": there is no window left for "fresh" to slip in before the sweep starts,
|
||||
// or to win the lock ahead of it — see the case-2 note below. This also means Case 1 (the sweep
|
||||
// is already inside its loop, holding the lock, when "fresh" tries to register) is now
|
||||
// guaranteed by construction rather than merely likely under a fixed sleep.
|
||||
assertTrue(sweepIsHoldingTheLock.await(20, TimeUnit.SECONDS),
|
||||
"the sweep never failed a single stale publish — it may not have started");
|
||||
|
||||
// Case 2 ("fresh" wins publishChannelLock before the sweep even starts, so it genuinely
|
||||
// published on the stale channel and the sweep correctly fails it) is impossible by
|
||||
// construction in this test: freshThread.start() below is reached only after
|
||||
// sweepIsHoldingTheLock has counted down, which can only happen once
|
||||
// failPendingPublishesOnRecovery() is already running and has already failed a stale entry.
|
||||
// There is no code path that lets "fresh" start before the sweep starts. That case is real
|
||||
// and correct production behaviour (see AmqpReplyInbox#failPendingPublishesOnRecovery's
|
||||
// javadoc), it is just not reachable from this deterministic ordering, so it does not need a
|
||||
// separate assertion here.
|
||||
|
||||
// This is the exact interleaving CB-528's follow-up describes: "the still-running recovery
|
||||
// sweep" racing a publish that registers while it is mid-flight.
|
||||
@@ -109,8 +133,8 @@ class AmqpReplyInboxRecoveryRaceTest {
|
||||
// Simulate the broker's real confirm for "fresh" now that the sweep is done, so a correct
|
||||
// implementation's publish() returns normally instead of idling out CONFIRM_TIMEOUT_MS. Poll
|
||||
// for the registration rather than checking once: freshThread may still be contending for
|
||||
// publishChannelLock (behind the 20,000 stale threads' own lock acquisitions) even though the
|
||||
// sweep itself has already finished.
|
||||
// publishChannelLock (behind the sweep's own hold on it, and possibly other stale threads
|
||||
// still unwinding) even though the sweep itself has already finished.
|
||||
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(9);
|
||||
int idx = -1;
|
||||
while (idx < 0 && System.nanoTime() < deadline) {
|
||||
|
||||
@@ -714,6 +714,47 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
// --- CB-582: bridge_status pendingAsk() ------------------------------------------------------
|
||||
|
||||
@Test
|
||||
void pendingAskReturnsNullWhenNoQuestionIsOpen() throws Exception {
|
||||
assertNull(messages.pendingAsk(T), "no async ticket at all -> no pending ask");
|
||||
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
awaitWaiting();
|
||||
assertNull(messages.pendingAsk(T), "a plain pending delegation is not a question");
|
||||
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||
assertNull(messages.pendingAsk(T), "a finished ticket carries no open question either");
|
||||
}
|
||||
|
||||
@Test
|
||||
void pendingAskReturnsTheOpenQuestionForAnAsyncTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
MessageService.PendingAsk pending = messages.pendingAsk(T);
|
||||
assertNotNull(pending, "bridge_status should see the open question");
|
||||
assertEquals(ticket, pending.ticket());
|
||||
assertEquals("which config file?", pending.question());
|
||||
assertEquals(asking.turnId(), pending.turnId());
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
assertNull(messages.pendingAsk(T), "an answered question is no longer pending");
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
// --- CB-588: async ticket terminal nudges ---------------------------------------------------
|
||||
//
|
||||
// MessageService.reply's rendezvous fast path is exactly what an async ticket always takes
|
||||
@@ -833,6 +874,110 @@ class MessageServiceTest {
|
||||
}
|
||||
}
|
||||
|
||||
// --- CB-582: bridge_ask question-open nudges --------------------------------------------------
|
||||
|
||||
@Test
|
||||
void anAsyncTicketThatPausesOnAQuestionNudgesTheLeadWithNoPriorPollCall() throws Exception {
|
||||
try (var wiring = wireWithPushLoop(5, 50)) {
|
||||
String ticket = wiring.service().sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync(
|
||||
() -> wiring.service().ask(T, "which config file?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.ASKING);
|
||||
|
||||
awaitNudge(wiring.leadHerdr());
|
||||
String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
|
||||
assertTrue(nudge.contains(ticket), "the nudge should name the ticket: " + nudge);
|
||||
assertTrue(nudge.contains(asking.turnId()), "the nudge should name the turnId: " + nudge);
|
||||
assertTrue(nudge.contains("bridge_send(turnId="),
|
||||
"the nudge should name the exact answer call: " + nudge);
|
||||
assertTrue(nudge.contains("which config file?"), "the nudge should include the question: " + nudge);
|
||||
|
||||
// Clean up the still-open ask so the background thread does not linger past the test.
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
answer.get(5, TimeUnit.SECONDS);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void answeringAQuestionStopsFurtherNudgesAboutIt() throws Exception {
|
||||
try (var wiring = wireWithPushLoop(5, 50)) {
|
||||
String ticket = wiring.service().sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync(
|
||||
() -> wiring.service().ask(T, "which config file?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.ASKING);
|
||||
|
||||
awaitNudge(wiring.leadHerdr());
|
||||
long callsBeforeAnswer = wiring.leadHerdr().calls.stream()
|
||||
.filter(c -> c.method().equals("agent.prompt")).count();
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
answer.get(5, TimeUnit.SECONDS);
|
||||
|
||||
// Let several more ticks (and the ticket's own now-legitimate terminal nudge) fire —
|
||||
// none of them may still name the question's turnId, which is closed.
|
||||
Thread.sleep(300);
|
||||
boolean anyNamesClosedQuestion = wiring.leadHerdr().calls.stream()
|
||||
.filter(c -> c.method().equals("agent.prompt"))
|
||||
.skip(callsBeforeAnswer)
|
||||
.anyMatch(c -> c.params().toString().contains(asking.turnId()));
|
||||
assertFalse(anyNamesClosedQuestion,
|
||||
"no nudge sent after the answer may still name the now-closed turnId " + asking.turnId());
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void anAskThatLeavesByThrowingStillClosesItsQuestion() throws Exception {
|
||||
// CB-582 follow-up. ask() calls clearAsyncQuestion on three paths — no-waiter, timed out,
|
||||
// and (from answer()) answered — but it can also leave by *throwing*: an interrupt while
|
||||
// blocked on the answer, or an ExecutionException from the answer future. Those paths run
|
||||
// only the finally block, so before the fix the push loop kept the question pending for
|
||||
// good: named in every nudge until its own cap, then never removed from the map at all.
|
||||
PrimaryRegistry registry = new PrimaryRegistry(null);
|
||||
registry.recordDelegation(T, LEAD);
|
||||
var scheduler = java.util.concurrent.Executors.newSingleThreadScheduledExecutor();
|
||||
// A backoff far longer than the test: the schedule is started but no tick ever fires, so
|
||||
// decide() is read directly and nothing here depends on timing.
|
||||
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, new AgentControl(new FakeHerdr()), inbox,
|
||||
scheduler, 5, 60_000);
|
||||
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop);
|
||||
try {
|
||||
String ticket = service.sendAsync(T, "task that asks");
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
Thread asker = new Thread(() -> assertThrows(IllegalStateException.class,
|
||||
() -> service.ask(T, "which config file?", 30_000)));
|
||||
asker.start();
|
||||
awaitTicketPhaseOn(service, ticket, MessageService.Phase.ASKING);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, pushLoop.decide(LEAD, 0, 0, 0),
|
||||
"the open question should be the one thing keeping this lead's schedule alive");
|
||||
|
||||
asker.interrupt();
|
||||
asker.join(5000);
|
||||
assertFalse(asker.isAlive(), "the interrupted ask should have left ask() by throwing");
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.STOP, pushLoop.decide(LEAD, 0, 0, 0),
|
||||
"an ask that threw must still close its question, or the loop nudges about it for good");
|
||||
} finally {
|
||||
service.close();
|
||||
scheduler.shutdownNow();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFleetWithNoPushLoopConfiguredBehavesExactlyAsToday() throws Exception {
|
||||
// `messages` (the shared field) uses the no-pushLoop constructor — poll() must not throw,
|
||||
|
||||
@@ -42,6 +42,7 @@ class ReplyPushLoopTest {
|
||||
|
||||
private static final String PRIMARY = "term_primary";
|
||||
private static final String WORKER = "term_worker";
|
||||
private static final String WORKER2 = "term_worker2";
|
||||
private static final ObjectMapper MAPPER = new ObjectMapper();
|
||||
|
||||
private PrimaryRegistry registry;
|
||||
@@ -579,6 +580,209 @@ class ReplyPushLoopTest {
|
||||
"both the reply and the ticket source are at their own cap — must still stop");
|
||||
}
|
||||
|
||||
// --- CB-598: work arriving during a backoff must not read as stale backlog -------------------
|
||||
|
||||
@Test
|
||||
void aTargetArrivingDuringTheBackoffGetsNudgedDespiteAnAlreadyCappedSibling() {
|
||||
// The bug: reminder counts used to be a single counter per lead per source, carried
|
||||
// forward across scheduled ticks (scheduleNext(lead, count + 1, ...)) rather than tracked
|
||||
// per pending item. WORKER gets nudged once here, which — with cap=1 — exhausts the
|
||||
// shared reply-source counter for this lead. WORKER2 then queues a reply for the SAME
|
||||
// lead "during the backoff": while the schedule from WORKER's tick is still active, before
|
||||
// the next tick's own start-of-tick snapshot runs. At that next tick, the OLD code passed
|
||||
// the already-exhausted shared counter into decide() regardless of WORKER2 never having
|
||||
// been named in any nudge, and — because WORKER2 was already present in that tick's
|
||||
// "before" snapshot — stopOrRestart's race check (proven correct on its own elsewhere in
|
||||
// this file) does not save it either: it looks like ordinary stale backlog, not a race.
|
||||
// WORKER2 was then stranded forever with no live schedule and no nudge ever naming it.
|
||||
//
|
||||
// tick() is driven directly (package-private, same reasoning as stopOrRestart being
|
||||
// directly testable) so the exact interleaving is deterministic instead of racing the
|
||||
// scheduler thread over a real ~15s backoff.
|
||||
//
|
||||
// Before the fix, this test fails on the second assertEquals: rec.sendCount() stays at 1
|
||||
// (decide() returns STOP on the second tick(), so injectNudge is never called a second
|
||||
// time) and the "must still get one" assertion never even runs.
|
||||
int cap = 1;
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
inbox.own(WORKER2);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop(cap, 100_000); // huge backoff — nothing fires on its own; we drive tick()
|
||||
|
||||
loop.onReplyQueued(WORKER);
|
||||
loop.tick(PRIMARY); // first tick: nudges WORKER alone; WORKER's own count reaches the cap
|
||||
assertEquals(1, rec.sendCount(), "the first tick should nudge about WORKER");
|
||||
|
||||
// WORKER2 "arrives during the backoff": queued for the same lead while the schedule from
|
||||
// the tick above is still active (activeLeads still holds PRIMARY), before the next tick
|
||||
// (simulated below) takes its own start-of-tick snapshot.
|
||||
inbox.publish(WORKER2, "m2", "hello2");
|
||||
loop.onReplyQueued(WORKER2);
|
||||
|
||||
loop.tick(PRIMARY); // the tick that would fire once that backoff elapsed
|
||||
|
||||
assertEquals(2, rec.sendCount(),
|
||||
"WORKER2 was never named in any nudge yet and must still get one, even though "
|
||||
+ "WORKER's own reminder count is already at the cap");
|
||||
String secondNudge = rec.sentParams().get(1).getValue().toString();
|
||||
assertTrue(secondNudge.contains(WORKER2), "the never-named target must be named: " + secondNudge);
|
||||
|
||||
// Criterion #3: isActive() must reflect that this lead still had a live nudge to give —
|
||||
// the second tick took the INJECT branch, so the schedule stayed live rather than being
|
||||
// torn down under WORKER2.
|
||||
assertTrue(loop.isActive(), "the schedule must stay active after nudging the fresh target");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aTicketArrivingDuringTheBackoffGetsNudgedDespiteAnAlreadyCappedSibling() {
|
||||
// Mirrors the reply-side test above for the ticket source.
|
||||
int cap = 1;
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
var loop = loop(cap, 100_000);
|
||||
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
loop.tick(PRIMARY); // first tick: nudges task-1 alone; its count reaches the cap
|
||||
assertEquals(1, rec.sendCount(), "the first tick should nudge about task-1");
|
||||
|
||||
loop.onTicketTerminal("task-2", WORKER, false); // arrives during the backoff, same lead
|
||||
loop.tick(PRIMARY);
|
||||
|
||||
assertEquals(2, rec.sendCount(),
|
||||
"task-2 was never named in any nudge yet and must still get one, even though "
|
||||
+ "task-1's reminder count is already at the cap");
|
||||
String secondNudge = rec.sentParams().get(1).getValue().toString();
|
||||
assertTrue(secondNudge.contains("task-2"), "the never-named ticket must be named: " + secondNudge);
|
||||
assertTrue(loop.isActive(), "the schedule must stay active after nudging the fresh ticket");
|
||||
}
|
||||
|
||||
// --- CB-582: bridge_ask question-open nudges -------------------------------------------------
|
||||
|
||||
@Test
|
||||
void onQuestionOpenedWithNoKnownLeadNeverStartsASchedule() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 50);
|
||||
|
||||
loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||
|
||||
assertFalse(loop.isActive(), "no lead known means nothing to nudge yet");
|
||||
Thread.sleep(150);
|
||||
assertEquals(0, rec.sendCount(), "must not nudge when no lead is known to be waiting");
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideQuestionsWithNothingPendingIsStop() {
|
||||
agents = agentWithStatus("idle");
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideQuestionsAtCapIsStop() {
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = loop(2, 100_000);
|
||||
loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0, 2));
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideQuestionsUnderCapWithInjectableLeadIsInject() {
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = loop(5, 100_000);
|
||||
loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void onQuestionOpenedCausesExactlyOneNudgeNamingTheTicketAndTurnId() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
|
||||
loop(1, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config file?");
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one question nudge should have been sent");
|
||||
assertEquals(1, rec.sendCount());
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("task-1"), "nudge should name the ticket: " + nudge);
|
||||
assertTrue(nudge.contains("term_worker#1"), "nudge should name the turnId: " + nudge);
|
||||
assertTrue(nudge.contains("bridge_send(turnId="), "nudge should name the exact answer call: " + nudge);
|
||||
assertTrue(nudge.contains("which config file?"), "nudge should include the question text: " + nudge);
|
||||
}
|
||||
|
||||
@Test
|
||||
void questionClosedPreventsFurtherNudging() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
var loop = loop(1, 100);
|
||||
|
||||
loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||
loop.questionClosed("term_worker#1"); // answered/lapsed before the first tick fired
|
||||
|
||||
Thread.sleep(300); // let the scheduled tick run
|
||||
assertEquals(0, rec.sendCount(), "an already-closed question must never be nudged");
|
||||
}
|
||||
|
||||
@Test
|
||||
void questionNudgesSendUpToCapThenStop() throws Exception {
|
||||
int cap = 2;
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
rec.sendLatch = new CountDownLatch(cap);
|
||||
|
||||
loop(cap, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||
|
||||
assertTrue(rec.sendLatch.await(5, TimeUnit.SECONDS), cap + " question nudges should have fired");
|
||||
Thread.sleep(300);
|
||||
assertEquals(cap, rec.sendCount(), "exactly " + cap + " question nudges (cap=" + cap + ")");
|
||||
}
|
||||
|
||||
@Test
|
||||
void questionAndTicketForTheSameLeadCoalesceIntoOneSend() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
var loop = loop(1, 300); // backoff wide enough that both entry points land before the first tick
|
||||
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
loop.onQuestionOpened("task-2", WORKER, "term_worker#1", "which config?");
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one combined nudge should have been sent");
|
||||
Thread.sleep(300);
|
||||
assertEquals(1, rec.sendCount(),
|
||||
"a ticket and a question for the same lead must coalesce onto ONE schedule");
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("task-1"), "the combined nudge must still mention the ticket: " + nudge);
|
||||
assertTrue(nudge.contains("term_worker#1"), "the combined nudge must still mention the question: " + nudge);
|
||||
}
|
||||
|
||||
@Test
|
||||
void oneExhaustedQuestionSourceDoesNotBlockANudgeForTheOtherSources() {
|
||||
// Mirrors oneExhaustedSourceDoesNotBlockANudgeForTheOtherSource for the question source:
|
||||
// the question source is at its cap (2/2), but the ticket source has never been nudged
|
||||
// (0/2) — decide() must still INJECT so the ticket is not stranded.
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = loop(2, 100_000);
|
||||
loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||
loop.onTicketTerminal("task-2", WORKER, false);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0, 2),
|
||||
"the question source is exhausted (2/2), but the ticket source has never been "
|
||||
+ "nudged (0/2) — the lead must still be injected so the ticket is not lost");
|
||||
}
|
||||
|
||||
@Test
|
||||
void questionNudgeFormatIsCorrect() {
|
||||
String single = ReplyPushLoop.QUESTION_NUDGE_FORMAT.formatted(
|
||||
WORKER, "task-1", "term_worker#1", "which config?");
|
||||
assertTrue(single.contains("Worker term_worker"));
|
||||
assertTrue(single.contains("bridge_send(turnId=\"term_worker#1\""));
|
||||
assertTrue(single.contains("which config?"));
|
||||
|
||||
String multi = ReplyPushLoop.QUESTIONS_NUDGE_FORMAT.formatted(2, "task-1 (turnId=t1), task-2 (turnId=t2)");
|
||||
assertTrue(multi.contains("2 workers"));
|
||||
assertTrue(multi.contains("bridge_poll(ticket=...)"));
|
||||
}
|
||||
|
||||
// --- metrics (CB-512) ----------------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -484,6 +484,67 @@ class BridgedAppTest {
|
||||
assertTrue(mapper.readTree(req(port, "GET", "/sessions/term_a/status").body()).get("ready").asBoolean());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-582: a lead polling {@code GET /sessions/{id}/status} on its normal cadence — not the
|
||||
* ticket-scoped {@code /tasks/{ticket}} — must also see a worker's open async {@code bridge_ask}
|
||||
* question, since the reverse-rendezvous window it opened with is far shorter than that cadence.
|
||||
*/
|
||||
@Test
|
||||
void sessionStatusReportsAnOpenQuestionWhenTheWorkerIsMidAsk() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().agentStatus("idle"); // poller delivers the injection
|
||||
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||
|
||||
HttpResponse<String> accepted = postMessage(port, "{\"content\":\"do it\",\"wait\":false}");
|
||||
assertEquals(202, accepted.statusCode());
|
||||
String ticket = mapper.readTree(accepted.body()).get("ticket").asText();
|
||||
|
||||
Thread.sleep(200); // let the background async send open its rendezvous waiter
|
||||
|
||||
var ask = java.util.concurrent.CompletableFuture.supplyAsync(() -> {
|
||||
try {
|
||||
return postJson(port, "/sessions/term_a/ask",
|
||||
"{\"question\":\"which config file?\",\"timeoutMs\":5000}");
|
||||
} catch (Exception e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
});
|
||||
|
||||
JsonNode task;
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
do {
|
||||
task = mapper.readTree(req(port, "GET", "/tasks/" + ticket).body());
|
||||
if ("asking".equals(task.path("phase").asText())) break;
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(10);
|
||||
} while (System.currentTimeMillis() < deadline);
|
||||
assertEquals("asking", task.get("phase").asText());
|
||||
String turnId = task.get("turnId").asText();
|
||||
|
||||
HttpResponse<String> status = req(port, "GET", "/sessions/term_a/status");
|
||||
assertEquals(200, status.statusCode());
|
||||
JsonNode body = mapper.readTree(status.body());
|
||||
assertEquals("idle", body.get("status").asText(), "the live status must still be reported");
|
||||
assertEquals("which config file?", body.get("question").asText());
|
||||
assertEquals(turnId, body.get("turnId").asText());
|
||||
assertEquals(ticket, body.get("ticket").asText());
|
||||
|
||||
// Answer it — via the same /message route bridge_send uses, keyed by turnId — so the
|
||||
// background ask thread does not linger past the test.
|
||||
var answer = java.util.concurrent.CompletableFuture.supplyAsync(() -> {
|
||||
try {
|
||||
return postJson(port, "/sessions/term_a/message",
|
||||
"{\"content\":\"config.yaml\",\"turnId\":\"" + turnId + "\",\"timeoutMs\":4000}");
|
||||
} catch (Exception e) {
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
});
|
||||
HttpResponse<String> askResponse = ask.get(6, java.util.concurrent.TimeUnit.SECONDS);
|
||||
assertEquals(200, askResponse.statusCode());
|
||||
assertEquals("config.yaml", mapper.readTree(askResponse.body()).get("answer").asText());
|
||||
postJson(port, "/sessions/term_a/reply", "{\"content\":\"done\"}");
|
||||
answer.get(6, java.util.concurrent.TimeUnit.SECONDS);
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopWorkerInPanePlacementClosesOnlyThePane() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
@@ -422,6 +422,28 @@ class WorktreeSessionManagerTest {
|
||||
"acceptance criterion 6: a failed ticket's detail must carry the snapshot ref");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseNotifiesTheListenerWithTheAgentSessionId() {
|
||||
// CB-584 (issue #65 criterion 5): a failed ticket's detail must also name the agent
|
||||
// session, alongside worktree/branch/snapshot, so a lead can resume the conversation
|
||||
// rather than only re-dispatch a fresh member onto the same files.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
|
||||
.withDirty(true);
|
||||
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
|
||||
java.util.List<SessionManager.ReleaseDetail> released = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
sessions.onRelease(released::add);
|
||||
MemberSession s = sessions.acquire("ltms-local", MemberRole.DEV, null, "/caller/proj", null,
|
||||
new WorktreeRequest("cb-584-e", null), "cb-584-session", null);
|
||||
|
||||
assertNotNull(s.agentSessionId(), "a named session must mint an agent session id to assert on");
|
||||
sessions.release(s.paneId());
|
||||
|
||||
assertEquals(1, released.size());
|
||||
SessionManager.ReleaseDetail detail = released.getFirst();
|
||||
assertEquals(s.agentSessionId(), detail.agentSessionId());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFailingSnapshotStillPreservesTheWorktreeStopsThePaneAndNotifies() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
@@ -82,8 +82,25 @@
|
||||
<key>RunAtLoad</key>
|
||||
<true/>
|
||||
|
||||
<!-- Restart on crash, but not in a tight loop if the config is bad (bridged fails fast on a
|
||||
non-loopback bind without token auth — that is a config error, not a transient one). -->
|
||||
<!--
|
||||
CB-600 — read this before assuming ThrottleInterval bounds anything. It paces restarts to at
|
||||
most one per 10s; it does NOT cap how many times launchd retries. If bridged fails fast on
|
||||
every start — a bad bridged.yaml, for example auth.mode: token with the token env var unset,
|
||||
which throws in main() before the daemon ever binds a port — launchd restarts it forever,
|
||||
once every 10s, until a human intervenes. LaunchAgents have no "give up after N attempts"
|
||||
primitive, so this is not something a config change here can fix.
|
||||
|
||||
That loop stops only two ways: (1) `launchctl unload -w ~/Library/LaunchAgents/dev.ltms.bridged.plist`,
|
||||
or (2) the underlying cause gets fixed, so the process starts successfully and stays up (no
|
||||
more exits to restart). scripts/redeploy-bridged.sh does not add a third way — it does not
|
||||
make bridged self-disable on a config error, on purpose: a fail-fast exit path that
|
||||
sometimes decides "this is unrecoverable, stop trying" is one more thing that can misfire,
|
||||
and a wrongly self-disabled daemon needs the exact same manual `launchctl load -w` recovery
|
||||
this comment already names — so it buys nothing an operator watching for the crash loop
|
||||
doesn't already have, at the cost of a new way to be silently down. Watch for it with
|
||||
`launchctl list dev.ltms.bridged` (a high restart count) or by tailing bridged.out for the
|
||||
same startup error repeating every ~10s.
|
||||
-->
|
||||
<key>KeepAlive</key>
|
||||
<dict>
|
||||
<key>SuccessfulExit</key>
|
||||
|
||||
@@ -79,6 +79,53 @@ running_pid() { pgrep -f "$PATTERN" || true; }
|
||||
launchd_installed() { [ -f "$LAUNCHD_PLIST" ]; }
|
||||
launchd_loaded() { launchctl list "$LAUNCHD_LABEL" >/dev/null 2>&1; }
|
||||
|
||||
# CB-600: the script computes its own log path from where it sits on disk (REPO, above); the
|
||||
# plist hard-codes an absolute StandardOutPath. Nothing forced the two to agree — if this script
|
||||
# were ever run from a checkout other than the one the loaded plist names, launchd would start and
|
||||
# log the daemon correctly, while every check below (the fresh "bridged listening" line, the
|
||||
# ERROR-count scan) would read a different, empty or stale file and the script would report a
|
||||
# clean restart while the daemon crash-loops. Pure and side-effect-free besides `die`/`ok` — reads
|
||||
# the two paths, resolves them, compares — so it never touches launchd or the daemon and can be
|
||||
# exercised by sourcing this script (see the SOURCED guard below) without installing the agent.
|
||||
check_log_path_matches_plist() {
|
||||
local script_out="$1" plist_path="$2"
|
||||
local plist_out resolved_out resolved_plist_out
|
||||
# Checked by exit status, not by emptiness: on a missing file/key PlistBuddy exits nonzero but
|
||||
# still writes a message ("File Doesn't Exist, Will Create: ...") that command substitution
|
||||
# would happily capture as if it were the real value — testing only `-z` missed that case.
|
||||
if ! plist_out="$(/usr/libexec/PlistBuddy -c 'Print :StandardOutPath' "$plist_path" 2>/dev/null)" \
|
||||
|| [ -z "$plist_out" ]; then
|
||||
die "launchd agent is loaded but PlistBuddy could not read StandardOutPath from
|
||||
$plist_path
|
||||
— cannot verify the daemon logs where this script is about to look. Fix the plist before
|
||||
redeploying supervised."
|
||||
fi
|
||||
resolved_out="$(cd "$(dirname "$script_out")" 2>/dev/null && pwd -P)/$(basename "$script_out")" || true
|
||||
resolved_plist_out="$(cd "$(dirname "$plist_out")" 2>/dev/null && pwd -P)/$(basename "$plist_out")" || true
|
||||
if [ -z "$resolved_out" ] || [ -z "$resolved_plist_out" ] || [ "$resolved_out" != "$resolved_plist_out" ]; then
|
||||
die "log path mismatch — this script reads
|
||||
$script_out (resolved: ${resolved_out:-<directory does not exist>})
|
||||
but the loaded plist's StandardOutPath is
|
||||
$plist_out (resolved: ${resolved_plist_out:-<directory does not exist>})
|
||||
Under supervision the daemon writes to the PLIST's path, not necessarily this script's — every
|
||||
post-restart check below (the fresh 'bridged listening' line, the ERROR-count scan) would read
|
||||
the wrong file and could report a clean restart while the daemon crash-loops. Fix the mismatch
|
||||
(move this checkout to match the plist, or edit the plist's StandardOutPath/StandardErrorPath)
|
||||
before redeploying supervised."
|
||||
fi
|
||||
ok "log path check: script and plist agree ($resolved_out)"
|
||||
}
|
||||
|
||||
# CB-600: sourceable for testing. When this file is SOURCED (not executed) it stops here — nothing
|
||||
# below runs — so a test harness can `source` it to call check_log_path_matches_plist (or the
|
||||
# other pure helpers above) against a throwaway plist fixture without ever reaching the mutating
|
||||
# flow (build/stop/start) or touching the real daemon or launchd. On a normal `./redeploy-bridged.sh`
|
||||
# invocation `(return 0 2>/dev/null)` fails (return is illegal at top level of an executed script),
|
||||
# so this whole block is a no-op and every line below still runs exactly as before.
|
||||
if (return 0 2>/dev/null); then
|
||||
return 0
|
||||
fi
|
||||
|
||||
# ---------------------------------------------------------------- report state
|
||||
|
||||
say "current state"
|
||||
@@ -103,6 +150,9 @@ SUPERVISED=0
|
||||
if launchd_loaded; then
|
||||
SUPERVISED=1
|
||||
ok "launchd agent loaded ($LAUNCHD_LABEL) — launchd supervises this daemon"
|
||||
# CB-600: fail loudly here, before ANY other check runs, if this script and the loaded plist
|
||||
# would read different log files — every check after this point is worthless otherwise.
|
||||
check_log_path_matches_plist "$OUT" "$LAUNCHD_PLIST"
|
||||
else
|
||||
warn "launchd agent not loaded — this script is the only thing that will restart the daemon."
|
||||
fi
|
||||
@@ -222,7 +272,22 @@ say "start"
|
||||
if [ "$SUPERVISED" = 1 ]; then
|
||||
echo " supervision is ON: using 'launchctl load' so launchd starts and keeps supervising this"
|
||||
echo " process, instead of a manual nohup that launchd would know nothing about."
|
||||
launchctl load -w "$LAUNCHD_PLIST" || die "launchctl load failed"
|
||||
# CB-600: 'launchctl unload -w' above already persisted Disabled=true for this label. A load -w
|
||||
# that succeeds clears it; a load -w that FAILS leaves the agent both stopped and disabled — worse
|
||||
# than before this script ran, because a later reboot or login will not bring it back either. One
|
||||
# retry covers a transient race (e.g. launchd not yet fully done deregistering); if it still fails,
|
||||
# die with the exact recovery command rather than a bare "failed".
|
||||
if ! launchctl load -w "$LAUNCHD_PLIST" 2>/dev/null; then
|
||||
warn "launchctl load failed on the first attempt — retrying once after a short pause"
|
||||
sleep 2
|
||||
launchctl load -w "$LAUNCHD_PLIST" || die "launchctl load failed twice.
|
||||
The agent is now STOPPED and DISABLED — it will NOT come back on its own, not even after a
|
||||
reboot or login, because 'launchctl unload -w' above persisted Disabled=true and load -w
|
||||
never got the chance to clear it. Recover with:
|
||||
launchctl load -w \"$LAUNCHD_PLIST\"
|
||||
If that still fails, check 'launchctl list $LAUNCHD_LABEL', validate the plist with
|
||||
'plutil -lint \"$LAUNCHD_PLIST\"', and check $OUT before assuming a retry will succeed."
|
||||
fi
|
||||
else
|
||||
# Absolute jar path so `ps` names which checkout is running.
|
||||
( cd "$BRIDGED" && zsh -lc "nohup java -jar '$JAR' >> bridged.out 2>&1 &" )
|
||||
|
||||
Reference in New Issue
Block a user