Compare commits
8 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 78ca24dc3f | |||
| 2124e043ce | |||
| e4f3620acb | |||
| 032a59a34d | |||
| e689090024 | |||
| 1cc34888fd | |||
| 0331ecd5d3 | |||
| 831a918c30 |
@@ -81,9 +81,9 @@ below are the procedure — run them in order, every task, not only the big ones
|
||||
5. **Collect** — `bridge_poll{ticket}` → `bridge_ack{ticket, msgId}`. Answer a worker's `bridge_ask`
|
||||
with `bridge_send{turnId, content}` — **not** `sessionId`. A worker gone quiet is diagnosed with
|
||||
`bridge_status`, never by reading its terminal.
|
||||
6. **Verify yourself.** Re-run the build and the checks. A worker mounts only the bridge MCP and
|
||||
cannot run your other tooling, and a piped command (`… | tail`) hides failures behind a zero
|
||||
exit — never promote a worker's "clean" to a fact.
|
||||
6. **Verify yourself.** Re-run the build and the checks. A worker cannot run your IDE tooling, any
|
||||
forge tools it appears to have hold a blocked credential and fail, and a piped command
|
||||
(`… | tail`) hides failures behind a zero exit — never promote a worker's "clean" to a fact.
|
||||
7. **Review — fan out.** Spawn reviewers against the diff, one per dimension or per file, with
|
||||
`wait:false`. Never the implementer of the scope it reviews, and brief them from the diff — not
|
||||
from the implementer's rationale, which carries its own blind spot. Dispatch each PR's reviewers
|
||||
@@ -155,9 +155,13 @@ you.
|
||||
without replying, the bridge scrapes your pane, and it can return only the last 4000 characters.
|
||||
A clipped scrape is marked as partial, but the missing text is gone — your report reaches the
|
||||
lead with its end cut off.
|
||||
5. **Report honestly.** State only what you actually ran and its real output, including failures.
|
||||
You mount **only** the bridge MCP — the primary's other servers (IDE, forge, docs) are not yours,
|
||||
so never claim the result of a check you had no way to run.
|
||||
5. **Report honestly.** State only what you actually ran and its real output, including failures,
|
||||
and never claim the result of a check you had no way to run. **Measure your own tools; do not
|
||||
assume them.** What you mount depends on your backend: an opencode member gets the bridge and
|
||||
nothing else, while a Claude Code member also inherits the operator's user-scope MCP servers,
|
||||
which the bridge never chose for you. Two rules follow. The primary's IDE tooling is still not
|
||||
yours, whatever you see. And **a mounted tool is not a working tool** — the forge server you may
|
||||
find there holds a deliberately blocked credential and fails every call, by design.
|
||||
6. **Never merge.** Stage files explicitly — never `git add -A` — and leave alone anything the
|
||||
project marks as not-yours-to-commit.
|
||||
|
||||
|
||||
@@ -781,10 +781,43 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* PATH} seeding already depends on the overlay reliably replacing an inherited value (see its
|
||||
* javadoc), and that is only demonstrated for a non-blank value, so this reuses the same,
|
||||
* proven-reliable shape rather than the unverified one.
|
||||
*
|
||||
* <p><b>MEASURED ON A LIVE PANE, 2026-08-15: this sentinel alone does NOT hold.</b> The overlay
|
||||
* itself works — {@code GITEA_TOKEN} is injected here, is exported by no shell file, and does
|
||||
* reach the pane. The sentinel loses one step later. A herdr pane runs a <em>login</em> shell,
|
||||
* {@code ~/.zprofile} sources {@code ${SHARED_ENV}/tools/secrets.sh}, and that file does a plain
|
||||
* unconditional {@code export GITEA_ACCESS_TOKEN=...}. A login shell overwrites a value already
|
||||
* in the environment, so the real admin token is put back over this sentinel before the member
|
||||
* process ever starts. That defeat applies to <em>every</em> name {@code secrets.sh} exports,
|
||||
* and no launcher-side overlay can win against it.
|
||||
*
|
||||
* <p>So this constant is not the control on its own — {@link #MEMBER_MARKER} is the other half.
|
||||
* Keeping the sentinel is still worth it: it is correct for any peer kind whose pane does not
|
||||
* start a login shell, and it makes the intent explicit at the one place every adapter passes.
|
||||
*/
|
||||
private static final String BLOCKED_GITEA_ACCESS_TOKEN =
|
||||
"blocked-by-bridged-cb592-see-gitea-issue-77";
|
||||
|
||||
/**
|
||||
* CB-592: marks a pane as a bridged member so a shell startup file can decline to export
|
||||
* operator-only credentials into it (gitea issue #77).
|
||||
*
|
||||
* <p>This name is deliberately one that {@code secrets.sh} never exports, which is exactly why
|
||||
* it survives the login shell that wipes {@link #BLOCKED_GITEA_ACCESS_TOKEN}. The mechanism is
|
||||
* measured, not assumed: {@code GITEA_TOKEN} is injected the same way, is absent from a login
|
||||
* shell of its own, and was observed set inside a live member pane.
|
||||
*
|
||||
* <p>It is a no-op until the operator guards the export, which is a one-line change in a file
|
||||
* this repo does not own and must not edit unasked:
|
||||
*
|
||||
* <pre>{@code
|
||||
* [ -n "${BRIDGED_MEMBER:-}" ] || export GITEA_ACCESS_TOKEN=...
|
||||
* }</pre>
|
||||
*
|
||||
* <p>Setting the marker now costs nothing and means that edit is the whole remaining fix.
|
||||
*/
|
||||
static final String MEMBER_MARKER = "BRIDGED_MEMBER";
|
||||
|
||||
/**
|
||||
* Seed a worker's environment (CB-511): the daemon's own {@code PATH}, then the profile's
|
||||
* {@code env:} entries, then the CB-592 admin-token shadow.
|
||||
@@ -802,10 +835,11 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
* dev.ltms.bridged.guard.SubscriptionGuard}, which is checked against the profile's
|
||||
* {@code baseUrl} and nothing else.
|
||||
*
|
||||
* <p>The CB-592 shadow is put in <em>last</em>, after the profile's own {@code env:}, so no
|
||||
* profile — present or future — can restore the admin token by naming it in config. This is
|
||||
* the one place the shadow is applied: every {@code buildLaunch} in every adapter calls this
|
||||
* first, so a new profile, and a peer kind not yet written, gets it for free.
|
||||
* <p>The CB-592 shadow and marker are put in <em>last</em>, after the profile's own
|
||||
* {@code env:}, so no profile — present or future — can restore the admin token, or hide that
|
||||
* the pane is a member, by naming either in config. This is the one place both are applied:
|
||||
* every {@code buildLaunch} in every adapter calls this first, so a new profile, and a peer
|
||||
* kind not yet written, gets them for free.
|
||||
*/
|
||||
protected Map<String, String> baseEnv(BridgedConfig.Profile cfg) {
|
||||
Map<String, String> workerEnv = new LinkedHashMap<>();
|
||||
@@ -817,6 +851,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
workerEnv.putAll(cfg.env());
|
||||
}
|
||||
workerEnv.put("GITEA_ACCESS_TOKEN", BLOCKED_GITEA_ACCESS_TOKEN);
|
||||
workerEnv.put(MEMBER_MARKER, "1");
|
||||
return workerEnv;
|
||||
}
|
||||
|
||||
|
||||
@@ -8,6 +8,8 @@ import dev.ltms.bridged.metrics.Metrics;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
@@ -16,35 +18,41 @@ import java.util.concurrent.TimeUnit;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
/**
|
||||
* Mechanism (b) of CB-307: a dedicated, status-gated push loop that nudges the primary's own
|
||||
* herdr pane when a worker reply lands with no live {@code bridge_send} to resolve it.
|
||||
* 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
|
||||
* async delegation ticket ({@code bridge_send(wait:false)}) that reached a terminal phase
|
||||
* (CB-588).
|
||||
*
|
||||
* <p>The loop is triggered by {@link #onReplyQueued(String)} (called from
|
||||
* {@link MessageService#reply} after the durable inbox publish). It checks four conditions
|
||||
* at each tick via {@link #decide(String, int)}, then either injects a drain nudge,
|
||||
* waits for the primary to become injectable, or stops reminding.
|
||||
* <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}'
|
||||
* compare-and-set), so at most one {@code agents.send} to that lead's pane is ever in flight.
|
||||
*
|
||||
* <p>Bounded: at most {@link #maxReminders} nudges per target, with a configurable backoff
|
||||
* between them. The reply is never lost — the durable inbox is the backstop.
|
||||
*
|
||||
* <p><strong>CB-588 ticket nudges.</strong> {@link #onTicketTerminal(String, String, boolean)} is a
|
||||
* second, independent entry point for an async delegation ticket ({@code bridge_send(wait:false)})
|
||||
* reaching a terminal phase. That path takes the rendezvous fast path in {@link MessageService#reply}
|
||||
* and never reaches {@link #onReplyQueued}, so without this the lead's own charter — prefer
|
||||
* {@code wait:false} for anything non-trivial — was exactly the mode this loop failed to cover. It
|
||||
* reuses the same status gating, bounded/backoff reminders, and metrics, keyed by the nudge-receiving
|
||||
* lead terminal rather than the worker target so several tickets finishing together coalesce into one
|
||||
* nudge. The two entry points do not interact: {@link #onReplyQueued} / {@link #decide} / their nudge
|
||||
* text and bound are unchanged.
|
||||
* <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)}). 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 the shared reminder cap
|
||||
* ({@link #maxReminders}) is reached, whichever the durable inbox / pending-ticket set doesn't
|
||||
* already answer via {@code STOP}.
|
||||
*/
|
||||
public final class ReplyPushLoop {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(ReplyPushLoop.class);
|
||||
static final String NUDGE_FORMAT = "Worker %s returned a reply — run bridge_poll(target=%s) to collect it";
|
||||
/** CB-588: singular form, one uncollected ticket. */
|
||||
/** Coalesced form, several uncollected replies for the same lead. */
|
||||
static final String REPLIES_NUDGE_FORMAT =
|
||||
"%d workers returned replies — run bridge_poll(target=...) for each to collect them: %s";
|
||||
/** Singular form, one uncollected ticket. */
|
||||
static final String TICKET_NUDGE_FORMAT =
|
||||
"Ticket %s finished%s — run bridge_poll(ticket=%s) to collect it";
|
||||
/** CB-588: coalesced form, several uncollected tickets for the same lead. */
|
||||
/** 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";
|
||||
|
||||
@@ -56,11 +64,11 @@ public final class ReplyPushLoop {
|
||||
private final long backoffMs;
|
||||
private final Metrics metrics; // CB-512: nullable — no registry in unit tests
|
||||
|
||||
/** Track targets that have an active schedule. */
|
||||
private final ConcurrentHashMap<String, Boolean> activeTargets = new ConcurrentHashMap<>();
|
||||
/** CB-588: tickets that have gone terminal but not yet been polled, keyed by ticket. */
|
||||
/** Worker targets with a reply queued, and the lead to nudge about it, keyed by target. */
|
||||
private final ConcurrentHashMap<String, String> 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-588: leads with an active ticket-reminder schedule. */
|
||||
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets). */
|
||||
private final ConcurrentHashMap<String, Boolean> activeLeads = new ConcurrentHashMap<>();
|
||||
|
||||
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||
@@ -89,178 +97,67 @@ public final class ReplyPushLoop {
|
||||
}
|
||||
}
|
||||
|
||||
// --- decision logic (package-private for unit-testing) -------------------------------------
|
||||
// --- pending-work lookups (package-private for unit-testing) -------------------------------
|
||||
|
||||
/** The action the loop should take for a target at the given reminder count. */
|
||||
/** The action the loop should take for a lead at the given reminder count. */
|
||||
enum Action { INJECT, WAIT_BUSY, STOP }
|
||||
|
||||
/**
|
||||
* Pure decision function: examine the current state and return what the loop should do.
|
||||
*
|
||||
* @param target the worker session (target terminal id)
|
||||
* @param reminderCount how many nudges have been sent so far for this target
|
||||
* @return the action the caller should take
|
||||
* Reply targets still pending for {@code lead} — registered via {@link #onReplyQueued} and
|
||||
* whose inbox still holds an unacked message. A target whose inbox has since drained (acked,
|
||||
* or collected via a live {@code bridge_send} rendezvous instead) is dropped from
|
||||
* {@link #pendingReplies} here rather than lingering forever; there is no explicit "reply
|
||||
* collected" callback the way {@link #ticketCollected} exists for tickets, so the inbox itself
|
||||
* is the only signal.
|
||||
*/
|
||||
Action decide(String target, int reminderCount) {
|
||||
// CB-532: the destination is per-delegation — the lead that sent this worker its work, not
|
||||
// "the primary". With two leads orchestrating one fleet the singular question has no right
|
||||
// answer, and answering it anyway interrupted whichever lead happened to call bridge_send
|
||||
// first with results it never asked for.
|
||||
var nudgeTarget = primaryRegistry.nudgeTargetFor(target);
|
||||
if (nudgeTarget.isEmpty()) {
|
||||
log.debug("push: no lead is known to be waiting on {}, stopping reminder", target);
|
||||
return Action.STOP;
|
||||
}
|
||||
if (inbox.peek(target).isEmpty()) {
|
||||
log.debug("push: inbox empty for {}, stopping reminder", target);
|
||||
return Action.STOP;
|
||||
}
|
||||
if (reminderCount >= maxReminders) {
|
||||
log.debug("push: reminder cap ({}) reached for {}, stopping", maxReminders, target);
|
||||
countNudge("exhausted");
|
||||
return Action.STOP;
|
||||
}
|
||||
String leadTerminal = nudgeTarget.get();
|
||||
AgentStatus status;
|
||||
try {
|
||||
status = agents.status(leadTerminal);
|
||||
} catch (RuntimeException e) {
|
||||
log.debug("push: status check failed for lead {}, will retry", leadTerminal, e);
|
||||
return Action.WAIT_BUSY;
|
||||
}
|
||||
if (status.injectable()) {
|
||||
return Action.INJECT;
|
||||
}
|
||||
log.debug("push: lead {} is {} (not injectable), waiting", leadTerminal, status);
|
||||
return Action.WAIT_BUSY;
|
||||
}
|
||||
|
||||
// --- public entrypoint ---------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Called when a reply is queued for {@code target}. Idempotent per target: a second call while
|
||||
* a schedule is active is a no-op. The schedule nudges the primary, then schedules a follow-up
|
||||
* check (reminder on backoff, or re-check on WAIT_BUSY), until the inbox is empty or the cap
|
||||
* is reached.
|
||||
*/
|
||||
public void onReplyQueued(String target) {
|
||||
if (activeTargets.putIfAbsent(target, Boolean.TRUE) != null) {
|
||||
log.debug("push: already active for {}, ignoring duplicate trigger", target);
|
||||
return; // already scheduled
|
||||
}
|
||||
log.debug("push: starting reminder loop for {}", target);
|
||||
scheduleNext(target, 0);
|
||||
}
|
||||
|
||||
/** Execute one loop tick — called on the scheduler thread. */
|
||||
private void tick(String target, int reminderCount) {
|
||||
var action = decide(target, reminderCount);
|
||||
switch (action) {
|
||||
case INJECT -> {
|
||||
injectNudge(target, reminderCount);
|
||||
scheduleNext(target, reminderCount + 1);
|
||||
}
|
||||
// Re-check after the configured backoff; the primary may become injectable soon.
|
||||
case WAIT_BUSY -> scheduleNext(target, reminderCount);
|
||||
case STOP -> {
|
||||
activeTargets.remove(target);
|
||||
log.debug("push: reminder loop ended for {}", target);
|
||||
private Set<String> pendingReplyTargetsFor(String lead) {
|
||||
Set<String> result = new HashSet<>();
|
||||
for (var entry : pendingReplies.entrySet()) {
|
||||
String target = entry.getKey();
|
||||
String owningLead = entry.getValue();
|
||||
if (!lead.equals(owningLead)) continue;
|
||||
if (inbox.peek(target).isEmpty()) {
|
||||
pendingReplies.remove(target, owningLead);
|
||||
continue;
|
||||
}
|
||||
result.add(target);
|
||||
}
|
||||
return result;
|
||||
}
|
||||
|
||||
/** Send the nudge and log the event. */
|
||||
private void injectNudge(String target, int reminderCount) {
|
||||
// Re-read rather than threading it down from decide(): the delegating lead can change
|
||||
// between the decision and the injection, and the nudge should follow the current one.
|
||||
var lead = primaryRegistry.nudgeTargetFor(target);
|
||||
if (lead.isEmpty()) {
|
||||
log.debug("push: lead for {} disappeared before the nudge could be sent", target);
|
||||
return;
|
||||
}
|
||||
String leadTerminal = lead.get();
|
||||
String nudge = NUDGE_FORMAT.formatted(target, target);
|
||||
try {
|
||||
agents.send(leadTerminal, nudge);
|
||||
log.debug("push: nudge {}/{} sent to lead {} for target {}",
|
||||
reminderCount + 1, maxReminders, leadTerminal, target);
|
||||
countNudge("delivered");
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("push: failed to nudge lead {} for target {} (reminder {}/{}): {}",
|
||||
leadTerminal, target, reminderCount + 1, maxReminders, e.toString());
|
||||
}
|
||||
}
|
||||
|
||||
/** Schedule the next tick on the scheduler thread pool. */
|
||||
private void scheduleNext(String target, int nextReminderCount) {
|
||||
scheduler.schedule(() -> tick(target, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS);
|
||||
}
|
||||
|
||||
// --- CB-588: async ticket terminal nudges ---------------------------------------------------
|
||||
|
||||
/** A ticket awaiting collection: which lead to nudge, and whether it ended in failure. */
|
||||
private record PendingTicket(String ticket, String lead, boolean failed) {
|
||||
}
|
||||
|
||||
/**
|
||||
* Called when an async delegation ticket ({@code bridge_send(wait:false)}, CB-107) reaches a
|
||||
* terminal phase — DONE or a failure. Unlike {@link #onReplyQueued}, which nudges about the
|
||||
* durable-inbox no-waiter path, this covers the path {@code MessageService.reply} takes when a
|
||||
* fire-and-poll send's own rendezvous waiter resolves the reply directly: that path returns
|
||||
* before {@link #onReplyQueued} is ever called, so without this entry point a ticket finishing
|
||||
* that way never nudged anyone (CB-588 / gitea #72).
|
||||
*
|
||||
* <p>Idempotent per lead: several tickets going terminal for the same lead while its schedule is
|
||||
* already active coalesce onto that schedule's next tick rather than firing a nudge each.
|
||||
*
|
||||
* @param ticket the ticket to nudge about
|
||||
* @param target the worker session the ticket was sent to — resolves which lead delegated it
|
||||
* @param failed whether the ticket ended in a failure phase rather than {@code DONE}
|
||||
*/
|
||||
public void onTicketTerminal(String ticket, String target, boolean failed) {
|
||||
var lead = primaryRegistry.nudgeTargetFor(target);
|
||||
if (lead.isEmpty()) {
|
||||
log.debug("push: no lead is known to be waiting on ticket {} (target {}), skipping nudge",
|
||||
ticket, target);
|
||||
return;
|
||||
}
|
||||
pendingTickets.put(ticket, new PendingTicket(ticket, lead.get(), failed));
|
||||
if (activeLeads.putIfAbsent(lead.get(), Boolean.TRUE) != null) {
|
||||
log.debug("push: ticket reminder loop already active for lead {}, {} coalesced in",
|
||||
lead.get(), ticket);
|
||||
return;
|
||||
}
|
||||
log.debug("push: starting ticket reminder loop for lead {}", lead.get());
|
||||
scheduleTicketTick(lead.get(), 0);
|
||||
}
|
||||
|
||||
/**
|
||||
* Called when a ticket's terminal state has been collected via {@code bridge_poll}. Removes it
|
||||
* from the pending set so a scheduled tick — and any nudge it sends — never names a ticket the
|
||||
* lead already has (CB-588 acceptance #5). A ticket that was never pending (unknown ticket, or
|
||||
* one nudged with no push loop configured) is a no-op.
|
||||
*/
|
||||
public void ticketCollected(String ticket) {
|
||||
pendingTickets.remove(ticket);
|
||||
}
|
||||
|
||||
/** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */
|
||||
private List<PendingTicket> pendingFor(String lead) {
|
||||
private List<PendingTicket> pendingTicketsFor(String lead) {
|
||||
return pendingTickets.values().stream().filter(t -> lead.equals(t.lead())).toList();
|
||||
}
|
||||
|
||||
/** Ticket ids still pending for {@code lead} — a plain snapshot for race comparison. */
|
||||
private Set<String> pendingTicketIdsFor(String lead) {
|
||||
return pendingTicketsFor(lead).stream().map(PendingTicket::ticket)
|
||||
.collect(Collectors.toUnmodifiableSet());
|
||||
}
|
||||
|
||||
/**
|
||||
* Pure decision function for ticket nudges, mirroring {@link #decide(String, int)} but keyed by
|
||||
* the nudge-receiving lead terminal rather than the worker session — several tickets from
|
||||
* different workers delegated by the same lead coalesce onto it.
|
||||
* Pure decision function: examine everything pending for {@code lead} — reply targets and
|
||||
* tickets alike — and return what the loop should do.
|
||||
*
|
||||
* @param lead the lead terminal to nudge
|
||||
* @param reminderCount how many nudges have been sent so far for this lead (shared across
|
||||
* both reply and ticket work — CB-590 collapses the reminder cap onto
|
||||
* one counter per lead, so alternating sources cannot outrun the bound)
|
||||
* @return the action the caller should take
|
||||
*/
|
||||
Action decideTickets(String lead, int reminderCount) {
|
||||
if (pendingFor(lead).isEmpty()) {
|
||||
log.debug("push: nothing pending for lead {}, stopping ticket reminder", lead);
|
||||
Action decide(String lead, int reminderCount) {
|
||||
boolean anyPending = !pendingReplyTargetsFor(lead).isEmpty() || !pendingTicketIdsFor(lead).isEmpty();
|
||||
if (!anyPending) {
|
||||
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
|
||||
return Action.STOP;
|
||||
}
|
||||
if (reminderCount >= maxReminders) {
|
||||
log.debug("push: ticket reminder cap ({}) reached for lead {}, stopping", maxReminders, lead);
|
||||
log.debug("push: reminder cap ({}) reached for lead {}, stopping", maxReminders, lead);
|
||||
countNudge("exhausted");
|
||||
return Action.STOP;
|
||||
}
|
||||
@@ -278,95 +175,185 @@ public final class ReplyPushLoop {
|
||||
return Action.WAIT_BUSY;
|
||||
}
|
||||
|
||||
/** Execute one ticket-loop tick — called on the scheduler thread. */
|
||||
private void ticketTick(String lead, int reminderCount) {
|
||||
Set<String> pendingBefore = pendingIdsFor(lead);
|
||||
var action = decideTickets(lead, reminderCount);
|
||||
switch (action) {
|
||||
case INJECT -> {
|
||||
injectTicketNudge(lead, reminderCount);
|
||||
scheduleTicketTick(lead, reminderCount + 1);
|
||||
}
|
||||
case WAIT_BUSY -> scheduleTicketTick(lead, reminderCount);
|
||||
case STOP -> stopOrRestartTicketLoop(lead, pendingBefore);
|
||||
}
|
||||
}
|
||||
// --- public entrypoints ----------------------------------------------------------------------
|
||||
|
||||
/** Ticket IDs pending for {@code lead} right now, as a plain snapshot for race comparison. */
|
||||
private Set<String> pendingIdsFor(String lead) {
|
||||
return pendingFor(lead).stream().map(PendingTicket::ticket).collect(Collectors.toUnmodifiableSet());
|
||||
/**
|
||||
* Called when a reply is queued for {@code target}. Resolves the lead delegating to
|
||||
* {@code target} (CB-532) and coalesces onto that lead's single reminder schedule — starting
|
||||
* one if none is active, joining an already-active one otherwise. A no-op if no lead is known
|
||||
* to be waiting on {@code target}: there is nobody to nudge yet, and the durable inbox is the
|
||||
* backstop until a lead is recorded.
|
||||
*/
|
||||
public void onReplyQueued(String target) {
|
||||
var lead = primaryRegistry.nudgeTargetFor(target);
|
||||
if (lead.isEmpty()) {
|
||||
log.debug("push: no lead is known to be waiting on {}, skipping reminder", target);
|
||||
return;
|
||||
}
|
||||
pendingReplies.put(target, lead.get());
|
||||
startOrCoalesce(lead.get());
|
||||
}
|
||||
|
||||
/**
|
||||
* Release {@code lead}'s active-schedule slot, then restart it only if a ticket landed that
|
||||
* {@code pendingBefore} — the snapshot taken just before this tick's decision — did not already
|
||||
* account for. {@code onTicketTerminal} reads {@code activeLeads} to decide whether to coalesce
|
||||
* onto an existing schedule or start one, so a ticket that lands between {@code decideTickets}
|
||||
* returning {@link Action#STOP} and this removal running sees the (soon-to-be-stale) slot as
|
||||
* occupied, coalesces onto a schedule that is about to die, and gets no nudge scheduled at all —
|
||||
* a lost nudge, the exact failure CB-588 exists to remove (found in review, gitea PR #73).
|
||||
* Called when an async delegation ticket ({@code bridge_send(wait:false)}, CB-107) reaches a
|
||||
* terminal phase — DONE or a failure. Unlike {@link #onReplyQueued}, which nudges about the
|
||||
* durable-inbox no-waiter path, this covers the path {@code MessageService.reply} takes when a
|
||||
* fire-and-poll send's own rendezvous waiter resolves the reply directly: that path returns
|
||||
* before {@link #onReplyQueued} is ever called, so without this entry point a ticket finishing
|
||||
* that way never nudged anyone (CB-588 / gitea #72).
|
||||
*
|
||||
* <p>Restarting on ANY non-empty {@code pendingFor(lead)} would be wrong: when STOP is reached
|
||||
* because the reminder cap was hit rather than the backlog draining, the same never-collected
|
||||
* ticket is expected to still be sitting there — that is the cap doing its job — and restarting
|
||||
* would nudge about it forever, defeating the bound (the original CB-307 bounded-reminder
|
||||
* guarantee, carried into CB-588 by acceptance criterion #7 — this exact regression showed up as
|
||||
* two existing tests failing once a naive "any pending ticket restarts" version of this fix went
|
||||
* in: {@code successfulTicketNudgeIncrementsDelivered} and {@code ticketNudgesSendUpToCapThenStop}).
|
||||
* Diffing the current pending set against {@code pendingBefore} tells the two cases apart: a
|
||||
* ticket present before this tick's decision is stale backlog, not a race; only a ticket absent
|
||||
* from {@code pendingBefore} can only have arrived during the decision-to-release window, which is
|
||||
* exactly the race this method closes.
|
||||
* <p>Resolves the delegating lead the same way {@link #onReplyQueued} does and coalesces onto
|
||||
* the same per-lead schedule (CB-590) — several tickets, or a ticket and a reply, finishing
|
||||
* for the same lead while its schedule is already active all ride the existing schedule's next
|
||||
* tick rather than firing a nudge each.
|
||||
*
|
||||
* @param ticket the ticket to nudge about
|
||||
* @param target the worker session the ticket was sent to — resolves which lead delegated it
|
||||
* @param failed whether the ticket ended in a failure phase rather than {@code DONE}
|
||||
*/
|
||||
public void onTicketTerminal(String ticket, String target, boolean failed) {
|
||||
var lead = primaryRegistry.nudgeTargetFor(target);
|
||||
if (lead.isEmpty()) {
|
||||
log.debug("push: no lead is known to be waiting on ticket {} (target {}), skipping nudge",
|
||||
ticket, target);
|
||||
return;
|
||||
}
|
||||
pendingTickets.put(ticket, new PendingTicket(ticket, lead.get(), failed));
|
||||
startOrCoalesce(lead.get());
|
||||
}
|
||||
|
||||
/**
|
||||
* Called when a ticket's terminal state has been collected via {@code bridge_poll}. Removes it
|
||||
* from the pending set so a scheduled tick — and any nudge it sends — never names a ticket the
|
||||
* lead already has. A ticket that was never pending (unknown ticket, or one nudged with no push
|
||||
* loop configured) is a no-op.
|
||||
*/
|
||||
public void ticketCollected(String ticket) {
|
||||
pendingTickets.remove(ticket);
|
||||
}
|
||||
|
||||
// --- the schedule ----------------------------------------------------------------------------
|
||||
|
||||
/** Start a reminder schedule for {@code lead}, or join the one already running. */
|
||||
private void startOrCoalesce(String lead) {
|
||||
if (activeLeads.putIfAbsent(lead, Boolean.TRUE) != null) {
|
||||
log.debug("push: reminder loop already active for lead {}, work coalesced in", lead);
|
||||
return;
|
||||
}
|
||||
log.debug("push: starting reminder loop for lead {}", lead);
|
||||
scheduleNext(lead, 0);
|
||||
}
|
||||
|
||||
/** Execute one loop tick — called on the scheduler thread. */
|
||||
private void tick(String lead, int reminderCount) {
|
||||
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
|
||||
Set<String> ticketsBefore = pendingTicketIdsFor(lead);
|
||||
var action = decide(lead, reminderCount);
|
||||
switch (action) {
|
||||
case INJECT -> {
|
||||
injectNudge(lead, reminderCount);
|
||||
scheduleNext(lead, reminderCount + 1);
|
||||
}
|
||||
// Re-check after the configured backoff; the lead may become injectable soon.
|
||||
case WAIT_BUSY -> scheduleNext(lead, reminderCount);
|
||||
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Release {@code lead}'s active-schedule slot, then restart it only if work landed that
|
||||
* {@code repliesBefore} / {@code ticketsBefore} — the snapshots taken just before this tick's
|
||||
* decision — did not already account for. {@link #onReplyQueued} / {@link #onTicketTerminal}
|
||||
* read {@link #activeLeads} to decide whether to coalesce onto an existing schedule or start
|
||||
* one, so work that lands between {@link #decide} returning {@link Action#STOP} and this
|
||||
* removal running sees the (soon-to-be-stale) slot as occupied, coalesces onto a schedule that
|
||||
* is about to die, and gets no nudge scheduled at all — a lost nudge, exactly what CB-588 (and
|
||||
* now CB-590) exist to remove (originally found in review, gitea PR #73, for the ticket-only
|
||||
* loop; carried forward here for the unified one).
|
||||
*
|
||||
* <p>Restarting on ANY non-empty pending set would be wrong: when STOP is reached because the
|
||||
* reminder cap was hit rather than the backlog draining, the same never-collected work is
|
||||
* expected to still be sitting there — that is the cap doing its job — and restarting would
|
||||
* nudge about it forever, defeating the bound. Diffing the current pending sets against the
|
||||
* "before" snapshots tells the two cases apart: an item present before this tick's decision is
|
||||
* stale backlog, not a race; only an item absent from the "before" snapshot can only have
|
||||
* arrived during the decision-to-release window, which is exactly the race this method closes.
|
||||
*
|
||||
* <p>Package-private so a test can drive the interleaving directly rather than trying to force a
|
||||
* genuine thread race: pass the exact {@code pendingBefore} snapshot a race requires (or does
|
||||
* not) and call this to prove the recheck responds correctly either way.
|
||||
* genuine thread race.
|
||||
*
|
||||
* <p>Terminates rather than spinning: this method restarts the schedule at most once per call, and
|
||||
* a fresh {@link #onTicketTerminal} racing the recheck below still terminates in one of two ways —
|
||||
* either it observes the slot already vacated (by the {@code activeLeads.remove} above, which
|
||||
* happens-before this recheck in program order) and claims it itself, or it lands first and this
|
||||
* recheck then observes its ticket in {@code pendingTickets} and reclaims the slot instead. Exactly
|
||||
* one side always wins; neither can miss the other, so this never loops on its own account.
|
||||
* <p>Terminates rather than spinning: this method restarts the schedule at most once per call,
|
||||
* and a fresh {@link #onReplyQueued} / {@link #onTicketTerminal} racing the recheck below still
|
||||
* terminates in one of two ways — either it observes the slot already vacated (by the
|
||||
* {@code activeLeads.remove} above, which happens-before this recheck in program order) and
|
||||
* claims it itself, or it lands first and this recheck then observes its work in
|
||||
* {@link #pendingReplies} / {@link #pendingTickets} and reclaims the slot instead. Exactly one
|
||||
* side always wins; neither can miss the other, so this never loops on its own account.
|
||||
*/
|
||||
void stopOrRestartTicketLoop(String lead, Set<String> pendingBefore) {
|
||||
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore) {
|
||||
activeLeads.remove(lead);
|
||||
boolean ticketRacedIn = pendingFor(lead).stream().anyMatch(t -> !pendingBefore.contains(t.ticket()));
|
||||
if (ticketRacedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
|
||||
log.debug("push: a ticket for lead {} raced the reminder loop's stop — restarting", lead);
|
||||
scheduleTicketTick(lead, 0);
|
||||
boolean racedIn = pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t))
|
||||
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.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);
|
||||
return;
|
||||
}
|
||||
log.debug("push: ticket reminder loop ended for lead {}", lead);
|
||||
log.debug("push: reminder loop ended for lead {}", lead);
|
||||
}
|
||||
|
||||
/** Send the coalesced ticket nudge and log the event. */
|
||||
private void injectTicketNudge(String lead, int reminderCount) {
|
||||
// Re-read rather than threading it down from decideTickets(): a ticket can be collected (or
|
||||
// another can arrive) between the decision and the injection.
|
||||
List<PendingTicket> pending = pendingFor(lead);
|
||||
if (pending.isEmpty()) {
|
||||
log.debug("push: pending tickets for lead {} drained before the nudge could be sent", lead);
|
||||
/** Send one combined nudge covering everything currently pending for {@code lead}. */
|
||||
private void injectNudge(String lead, int reminderCount) {
|
||||
// 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.
|
||||
Set<String> replyTargets = pendingReplyTargetsFor(lead);
|
||||
List<PendingTicket> tickets = pendingTicketsFor(lead);
|
||||
if (replyTargets.isEmpty() && tickets.isEmpty()) {
|
||||
log.debug("push: pending work for lead {} drained before the nudge could be sent", lead);
|
||||
return;
|
||||
}
|
||||
String nudge = formatTicketsNudge(pending);
|
||||
String nudge = formatNudge(replyTargets, tickets);
|
||||
try {
|
||||
agents.send(lead, nudge);
|
||||
log.debug("push: ticket nudge {}/{} sent to lead {} for {} ticket(s)",
|
||||
reminderCount + 1, maxReminders, lead, pending.size());
|
||||
log.debug("push: nudge {}/{} sent to lead {} ({} reply target(s), {} ticket(s))",
|
||||
reminderCount + 1, maxReminders, lead, replyTargets.size(), tickets.size());
|
||||
countNudge("delivered");
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("push: failed to nudge lead {} for {} ticket(s) (reminder {}/{}): {}",
|
||||
lead, pending.size(), reminderCount + 1, maxReminders, e.toString());
|
||||
log.warn("push: failed to nudge lead {} (reminder {}/{}): {}",
|
||||
lead, reminderCount + 1, maxReminders, e.toString());
|
||||
}
|
||||
}
|
||||
|
||||
/** Schedule the next ticket-loop tick on the scheduler thread pool. */
|
||||
private void scheduleTicketTick(String lead, int nextReminderCount) {
|
||||
scheduler.schedule(() -> ticketTick(lead, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS);
|
||||
/** Schedule the next tick on the scheduler thread pool. */
|
||||
private void scheduleNext(String lead, int nextReminderCount) {
|
||||
scheduler.schedule(() -> tick(lead, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS);
|
||||
}
|
||||
|
||||
/** Render one or several pending tickets as a single nudge line. */
|
||||
// --- nudge formatting ------------------------------------------------------------------------
|
||||
|
||||
/** Render everything pending for one lead as a single nudge line. */
|
||||
private static String formatNudge(Set<String> replyTargets, List<PendingTicket> tickets) {
|
||||
List<String> parts = new ArrayList<>();
|
||||
if (!replyTargets.isEmpty()) {
|
||||
parts.add(formatRepliesNudge(replyTargets));
|
||||
}
|
||||
if (!tickets.isEmpty()) {
|
||||
parts.add(formatTicketsNudge(tickets));
|
||||
}
|
||||
return String.join(" | ", parts);
|
||||
}
|
||||
|
||||
/** Render one or several pending reply targets. */
|
||||
private static String formatRepliesNudge(Set<String> targets) {
|
||||
if (targets.size() == 1) {
|
||||
String target = targets.iterator().next();
|
||||
return NUDGE_FORMAT.formatted(target, target);
|
||||
}
|
||||
String ids = String.join(", ", targets);
|
||||
return REPLIES_NUDGE_FORMAT.formatted(targets.size(), ids);
|
||||
}
|
||||
|
||||
/** Render one or several pending tickets. */
|
||||
private static String formatTicketsNudge(List<PendingTicket> pending) {
|
||||
if (pending.size() == 1) {
|
||||
PendingTicket t = pending.get(0);
|
||||
@@ -383,23 +370,22 @@ public final class ReplyPushLoop {
|
||||
// --- lifecycle -----------------------------------------------------------------------------
|
||||
|
||||
/**
|
||||
* Whether any reminder loop is currently active for some target (CB-551). The idle-lead heartbeat
|
||||
* uses this to stand aside: while the push loop is actively nudging the 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 #activeTargets} or
|
||||
* {@link #activeLeads} (CB-588 ticket nudges are a second source of pane injections the heartbeat
|
||||
* must equally stand aside for); the sets are bounded by what has been triggered, not by any
|
||||
* persistent state.
|
||||
* Whether any reminder loop is currently active for some lead (CB-551). The idle-lead heartbeat
|
||||
* 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.
|
||||
*/
|
||||
public boolean isActive() {
|
||||
return !activeTargets.isEmpty() || !activeLeads.isEmpty();
|
||||
return !activeLeads.isEmpty();
|
||||
}
|
||||
|
||||
/** Shut down the scheduler. Outstanding reminders are cancelled. */
|
||||
public void stop() {
|
||||
scheduler.shutdownNow();
|
||||
activeTargets.clear();
|
||||
activeLeads.clear();
|
||||
pendingReplies.clear();
|
||||
pendingTickets.clear();
|
||||
}
|
||||
|
||||
|
||||
@@ -685,6 +685,12 @@ class ClaudeCodeLauncherTest {
|
||||
* explicit (non-blank) GITEA_ACCESS_TOKEN to herdr on every spawn, whatever the profile is, so
|
||||
* a future baseEnv refactor cannot silently drop it and reopen the leak. Asserted against what
|
||||
* tab.create's params actually carry, not an internal map built in the test (gitea #77).
|
||||
*
|
||||
* <p>Scope, measured on a live pane 2026-08-15: this pins what the launcher SENDS, and that is
|
||||
* all it can pin. It does not prove the value survives, and it does not: the pane runs a login
|
||||
* shell, ~/.zprofile sources secrets.sh, and its unconditional `export GITEA_ACCESS_TOKEN=...`
|
||||
* puts the real token back over this sentinel. Closing that needs the operator to guard the
|
||||
* export on BRIDGED_MEMBER — see everySpawnMarksThePaneAsAMember below.
|
||||
*/
|
||||
@Test
|
||||
void everySpawnShadowsTheAdminGiteaAccessToken() {
|
||||
@@ -716,6 +722,39 @@ class ClaudeCodeLauncherTest {
|
||||
"a profile's own env: must not be able to smuggle the admin token back in");
|
||||
}
|
||||
|
||||
/**
|
||||
* The half of CB-592 that can actually survive the pane's login shell. BRIDGED_MEMBER is a name
|
||||
* secrets.sh never exports, so nothing overwrites it — measured: GITEA_TOKEN is injected the
|
||||
* same way, is absent from a login shell of its own, and was observed set inside a live member
|
||||
* pane. It lets the operator guard the admin export with
|
||||
* `[ -n "${BRIDGED_MEMBER:-}" ] || export GITEA_ACCESS_TOKEN=...`, which is the whole fix.
|
||||
* Pinned here so a refactor cannot drop the marker and quietly un-guard every member (#77).
|
||||
*/
|
||||
@Test
|
||||
void everySpawnMarksThePaneAsAMember() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
service(herdr, List.of("claude"), null).spawn();
|
||||
|
||||
assertEquals("1", startEnv(herdr).get("BRIDGED_MEMBER"),
|
||||
"every member pane must be marked, or a shell file cannot tell it apart from the operator's");
|
||||
}
|
||||
|
||||
/** A profile must not be able to hide that its pane is a member, for the same reason as above. */
|
||||
@Test
|
||||
void aProfileEnvEntryCannotClearTheMemberMarker() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
|
||||
List.of("claude"), "tab", "bridged-workers", "w #{n}", null, null, null, null, null,
|
||||
null, Map.of("BRIDGED_MEMBER", ""), null, null);
|
||||
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||
_ -> null).spawn();
|
||||
|
||||
assertEquals("1", startEnv(herdr).get("BRIDGED_MEMBER"),
|
||||
"a profile's own env: must not be able to unmark its pane");
|
||||
}
|
||||
|
||||
// ── CB-533: the model is pinned on the command line, not only in the environment ────────────
|
||||
|
||||
/** A launcher for a profile identical but for its {@code model:} — the only variable here. */
|
||||
|
||||
@@ -20,6 +20,7 @@ import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
@@ -27,6 +28,12 @@ import static org.junit.jupiter.api.Assertions.*;
|
||||
* Unit tests for {@link ReplyPushLoop}: decision logic, nudge injection, idempotency,
|
||||
* bounded reminders, and stop conditions.
|
||||
*
|
||||
* <p>CB-590 collapsed the CB-307 reply-nudge schedule and the CB-588 ticket-nudge schedule into
|
||||
* one schedule per lead ({@link ReplyPushLoop#decide}), so most tests below register pending work
|
||||
* through the public entry points ({@code onReplyQueued} / {@code onTicketTerminal}) before
|
||||
* exercising {@code decide} directly, mirroring how the two entry points now share one decision
|
||||
* function keyed by the lead terminal rather than by worker target.
|
||||
*
|
||||
* <p>Uses a {@link RecordingHerdrClient} that synchronizes access to its call list so the
|
||||
* scheduler thread and test thread never have memory ordering issues. The {@code decide()}
|
||||
* tests use a simple client with no concurrency concern.
|
||||
@@ -58,38 +65,50 @@ class ReplyPushLoopTest {
|
||||
// --- decide() logic ------------------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
void decideWithoutPrimaryIsStop() {
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = new ReplyPushLoop(
|
||||
new PrimaryRegistry(null), agents, inbox, scheduler, 5, 100);
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(WORKER, 0));
|
||||
void onReplyQueuedWithNoKnownLeadNeverStartsASchedule() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 50);
|
||||
|
||||
loop.onReplyQueued(WORKER); // no lead known -> never registered, never scheduled
|
||||
|
||||
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 decideWithEmptyInboxIsStop() {
|
||||
void decideWithNothingPendingIsStop() {
|
||||
agents = agentWithStatus("idle");
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideAtCapIsStop() {
|
||||
agents = agentWithStatus("idle");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100).decide(WORKER, 2));
|
||||
var loop = loop(2, 100);
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideUnderCapWithInjectablePrimaryIsInject() {
|
||||
agents = agentWithStatus("idle");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideUnderCapWithBlockedPrimaryIsInject() {
|
||||
agents = agentWithStatus("blocked");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0),
|
||||
"BLOCKED is injectable");
|
||||
}
|
||||
|
||||
@@ -97,7 +116,9 @@ class ReplyPushLoopTest {
|
||||
void decideUnderCapWithDonePrimaryIsInject() {
|
||||
agents = agentWithStatus("done");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0),
|
||||
"DONE is injectable");
|
||||
}
|
||||
|
||||
@@ -105,23 +126,29 @@ class ReplyPushLoopTest {
|
||||
void decideUnderCapWithBusyPrimaryIsWaitBusy() {
|
||||
agents = agentWithStatus("working");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideUnderCapWithUnknownPrimaryIsWaitBusy() {
|
||||
agents = agentWithStatus("unknown");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
void decideStopsAfterInboxIsEmptied() {
|
||||
agents = agentWithStatus("idle");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
|
||||
inbox.ack(WORKER, "m1");
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
// --- onReplyQueued integration -------------------------------------------------------------
|
||||
@@ -201,12 +228,19 @@ class ReplyPushLoopTest {
|
||||
assertTrue(nudge.contains("bridge_poll(target=term_worker)"));
|
||||
}
|
||||
|
||||
// --- CB-588: async ticket terminal nudges — decideTickets() logic --------------------------
|
||||
@Test
|
||||
void repliesNudgeFormatIsCorrect() {
|
||||
String multi = ReplyPushLoop.REPLIES_NUDGE_FORMAT.formatted(2, "term_worker1, term_worker2");
|
||||
assertTrue(multi.contains("2 workers"));
|
||||
assertTrue(multi.contains("bridge_poll(target=...)"));
|
||||
}
|
||||
|
||||
// --- CB-588: async ticket terminal nudges — decide() logic on tickets -----------------------
|
||||
|
||||
@Test
|
||||
void decideTicketsWithNothingPendingIsStop() {
|
||||
agents = agentWithStatus("idle");
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decideTickets(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -214,7 +248,7 @@ class ReplyPushLoopTest {
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = loop(2, 100_000);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -222,7 +256,7 @@ class ReplyPushLoopTest {
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = loop(5, 100_000);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decideTickets(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -230,7 +264,7 @@ class ReplyPushLoopTest {
|
||||
agents = agentWithStatus("working");
|
||||
var loop = loop(5, 100_000);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decideTickets(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -238,7 +272,7 @@ class ReplyPushLoopTest {
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 100_000);
|
||||
loop.onTicketTerminal("task-1", WORKER, false); // no lead known -> never registered as pending
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0));
|
||||
}
|
||||
|
||||
// --- CB-588: onTicketTerminal integration ---------------------------------------------------
|
||||
@@ -255,7 +289,7 @@ class ReplyPushLoopTest {
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("task-1"), "nudge should name the ticket");
|
||||
assertTrue(nudge.contains("bridge_poll(ticket="), "nudge should name the exact ticket-poll call");
|
||||
assertFalse(nudge.contains("bridge_poll(target="), "a ticket nudge must not tell the lead to run the target-poll call");
|
||||
assertFalse(nudge.contains("bridge_poll(target="), "a ticket-only nudge must not tell the lead to run the target-poll call");
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -346,11 +380,11 @@ class ReplyPushLoopTest {
|
||||
void aTicketStillPendingWhenTheLoopStopsIsNotStranded() {
|
||||
// Regression for the race a reviewer found in gitea PR #73: onTicketTerminal's
|
||||
// activeLeads.putIfAbsent can see the lead's slot as still occupied a moment before
|
||||
// decideTickets' STOP releases it, so the ticket coalesces onto a schedule that is about to
|
||||
// decide's STOP releases it, so the ticket coalesces onto a schedule that is about to
|
||||
// die and nothing ever nudges about it. Forcing that exact thread interleaving is not
|
||||
// reliable, so this drives stopOrRestartTicketLoop — the STOP path's own release-and-recheck —
|
||||
// reliable, so this drives stopOrRestart — the STOP path's own release-and-recheck —
|
||||
// directly, arranging the state it must not lose a ticket in: a ticket pending for the lead
|
||||
// that was NOT part of the pre-decision snapshot (pendingBefore=empty), standing in for one
|
||||
// that was NOT part of the pre-decision snapshot (ticketsBefore=empty), standing in for one
|
||||
// that races in during the decision-to-release window.
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
@@ -358,9 +392,9 @@ class ReplyPushLoopTest {
|
||||
|
||||
loop.onTicketTerminal("task-1", WORKER, false); // pendingTickets={task-1}; activeLeads={PRIMARY}
|
||||
|
||||
// Stand in for the scheduler thread reaching decideTickets==STOP for this lead — with nothing
|
||||
// Stand in for the scheduler thread reaching decide==STOP for this lead — with nothing
|
||||
// pending at decide time — while task-1 races in before the release below runs.
|
||||
loop.stopOrRestartTicketLoop(PRIMARY, Set.of());
|
||||
loop.stopOrRestart(PRIMARY, Set.of(), Set.of());
|
||||
|
||||
assertTrue(loop.isActive(), "a ticket that raced the loop's stop must reclaim the schedule "
|
||||
+ "slot, not be stranded with no schedule left to ever nudge about it");
|
||||
@@ -368,18 +402,19 @@ class ReplyPushLoopTest {
|
||||
|
||||
@Test
|
||||
void aStaleUncollectedTicketAtCapDoesNotRestartTheLoop() {
|
||||
// The other direction of the same fix: restarting on ANY non-empty pendingFor(lead) would be
|
||||
// The other direction of the same fix: restarting on ANY non-empty pending set would be
|
||||
// wrong. When STOP is reached because the reminder cap was hit, the same never-collected
|
||||
// ticket is expected to still be there — that is the cap doing its job (acceptance criterion
|
||||
// #7: nudges stay bounded). task-1 here was already accounted for at decide time (it is in
|
||||
// pendingBefore), so it must not restart the loop just because it is still sitting there.
|
||||
// #5: no spin / nudges stay bounded). task-1 here was already accounted for at decide time
|
||||
// (it is in ticketsBefore), so it must not restart the loop just because it is still sitting
|
||||
// there.
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
ReplyPushLoop loop = loop(1, 100_000);
|
||||
|
||||
loop.onTicketTerminal("task-1", WORKER, false); // pendingTickets={task-1}; activeLeads={PRIMARY}
|
||||
|
||||
loop.stopOrRestartTicketLoop(PRIMARY, Set.of("task-1"));
|
||||
loop.stopOrRestart(PRIMARY, Set.of(), Set.of("task-1"));
|
||||
|
||||
assertFalse(loop.isActive(), "a stale ticket already accounted for at decide time must not "
|
||||
+ "restart the loop — that would defeat the reminder cap");
|
||||
@@ -402,7 +437,52 @@ class ReplyPushLoopTest {
|
||||
void inboxNudgeStillUsesTheOriginalTargetPollCall() {
|
||||
String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER);
|
||||
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
|
||||
"CB-588 must not change the CB-307 inbox nudge's call shape");
|
||||
"CB-588/CB-590 must not change the CB-307 inbox nudge's call shape");
|
||||
}
|
||||
|
||||
// --- CB-590: one schedule per lead — no overlap, no lost nudges -----------------------------
|
||||
|
||||
@Test
|
||||
void replyAndTicketForTheSameLeadCoalesceIntoOneSendNeverOverlapping() throws Exception {
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop(1, 300); // backoff wide enough that both entry points land before the first tick
|
||||
|
||||
loop.onReplyQueued(WORKER);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one combined nudge should have been sent");
|
||||
Thread.sleep(300);
|
||||
assertEquals(1, rec.sendCount(),
|
||||
"a reply and a ticket for the same lead must coalesce onto ONE schedule — "
|
||||
+ "two nudge injections into the same lead pane must never overlap");
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
|
||||
"the combined nudge must still mention the reply: " + nudge);
|
||||
assertTrue(nudge.contains("task-1"), "the combined nudge must still mention the ticket: " + nudge);
|
||||
}
|
||||
|
||||
@Test
|
||||
void aReplyQueuedWhileTheLeadIsBusyIsNotLostWhenATicketArrivesToo() throws Exception {
|
||||
// The lead is busy for its first two status checks, then becomes injectable. A reply is
|
||||
// queued while busy; a ticket for the same lead arrives before the lead frees up. Neither
|
||||
// may be dropped — deferred is fine, lost is not (acceptance criterion #2).
|
||||
var rec = new BusyThenIdleHerdrClient(2);
|
||||
agents = new AgentControl(rec);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop(1, 50); // cap=1: WAIT_BUSY doesn't count against it, so exactly one send once injectable
|
||||
|
||||
loop.onReplyQueued(WORKER); // schedule starts, first tick(s) WAIT_BUSY
|
||||
loop.onTicketTerminal("task-1", WORKER, false); // coalesces onto the same waiting schedule
|
||||
|
||||
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS),
|
||||
"once the lead becomes injectable, the deferred work must still be nudged");
|
||||
Thread.sleep(200);
|
||||
assertEquals(1, rec.sendCount(), "exactly one nudge once injectable — reply and ticket coalesced");
|
||||
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"), "the reply must not be dropped: " + nudge);
|
||||
assertTrue(nudge.contains("task-1"), "the ticket must not be dropped: " + nudge);
|
||||
}
|
||||
|
||||
// --- metrics (CB-512) ----------------------------------------------------------------------
|
||||
@@ -430,8 +510,10 @@ class ReplyPushLoopTest {
|
||||
agents = agentWithStatus("idle");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
Metrics metrics = new Metrics();
|
||||
var loop = loop(2, 100, metrics);
|
||||
loop.onReplyQueued(WORKER);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100, metrics).decide(WORKER, 2));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
|
||||
|
||||
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
|
||||
"hitting the reminder cap must count as exhausted");
|
||||
@@ -459,7 +541,7 @@ class ReplyPushLoopTest {
|
||||
var loop = loop(2, 100_000, metrics);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
|
||||
|
||||
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
|
||||
"hitting the ticket reminder cap must count as exhausted");
|
||||
@@ -547,4 +629,49 @@ class ReplyPushLoopTest {
|
||||
private static RecordingHerdrClient recordingClient() {
|
||||
return new RecordingHerdrClient();
|
||||
}
|
||||
|
||||
/**
|
||||
* Thread-safe fake that reports {@code working} (not injectable) for its first
|
||||
* {@code busyChecks} status calls, then {@code idle} forever after — used to prove work queued
|
||||
* while the lead is busy is deferred, not dropped, once it becomes injectable.
|
||||
*/
|
||||
private static final class BusyThenIdleHerdrClient implements HerdrClient {
|
||||
private final List<Map.Entry<String, Object>> calls =
|
||||
Collections.synchronizedList(new ArrayList<>());
|
||||
private final AtomicInteger statusChecks = new AtomicInteger();
|
||||
private final int busyChecks;
|
||||
volatile CountDownLatch sendLatch = new CountDownLatch(1);
|
||||
|
||||
BusyThenIdleHerdrClient(int busyChecks) {
|
||||
this.busyChecks = busyChecks;
|
||||
}
|
||||
|
||||
@Override
|
||||
public JsonNode call(String method, Object params) {
|
||||
if ("agent.get".equals(method)) {
|
||||
String status = statusChecks.getAndIncrement() < busyChecks ? "working" : "idle";
|
||||
return MAPPER.createObjectNode()
|
||||
.set("agent", MAPPER.createObjectNode()
|
||||
.put("terminal_id", PRIMARY)
|
||||
.put("agent_status", status));
|
||||
}
|
||||
if ("agent.prompt".equals(method)) {
|
||||
calls.add(Map.entry(method, params));
|
||||
sendLatch.countDown();
|
||||
}
|
||||
return MAPPER.createObjectNode();
|
||||
}
|
||||
|
||||
long sendCount() {
|
||||
return calls.size();
|
||||
}
|
||||
|
||||
List<Map.Entry<String, Object>> sentParams() {
|
||||
return List.copyOf(calls);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,11 @@
|
||||
# CB-591 — move the fleet onto the LLM and MCP gateway
|
||||
|
||||
**Status:** plan, not started · **Upstream:** [systems/vms wiki → LLM and MCP Gateway](https://git.ltms.dev/systems/vms/wiki/LLM-and-MCP-Gateway)
|
||||
**Status: DONE — the fleet is on the gateway as of 2026-08-15.** `local` runs on `/anthropic` and
|
||||
`gx` on `/v1`, both at `weight: 100`; `local-direct` stays at `weight: 0` as the escape hatch. Getting
|
||||
here took a revert and two upstream fixes — see §7.1, which is the useful part of this document. One
|
||||
risk is **accepted rather than solved**: a stream cut by any mid-response timer arrives as HTTP 200
|
||||
with no terminator, and our third-party members cannot detect it (§7.2).
|
||||
· **Upstream:** [systems/vms wiki → LLM and MCP Gateway](https://git.ltms.dev/systems/vms/wiki/LLM-and-MCP-Gateway)
|
||||
· **Upstream issue:** [systems/vms#31](https://git.ltms.dev/systems/vms/issues/31)
|
||||
|
||||
The gateway went live on 2026-08-15 and replaced Bifrost. This plan says what that means for a
|
||||
@@ -154,12 +159,21 @@ speaks the Anthropic protocol, so `/anthropic` is both correct and the only safe
|
||||
This is the exact failure shape this repo keeps hitting: it compiles, it answers, it looks healthy,
|
||||
and a capability is quietly off. Treat it as a `silent-default` risk, not a config preference.
|
||||
|
||||
**Open question for 3b, stated as open.** The opencode profile must use the OpenAI surface, because
|
||||
that is the only thing an OpenAI-compatible provider can speak. vLLM serves that API natively, so I
|
||||
*expect* it to be a passthrough as well — but the wiki documents the translation trap only for the
|
||||
Anthropic path, and I have not checked how `reasoning_content` behaves through `/v1`. Verify it on
|
||||
first spawn (§7) rather than assuming. If reasoning is dropped there, that is a limitation of the
|
||||
opencode profile, not a reason to abandon it — opencode members do dev work, not deep reasoning.
|
||||
**Open question for 3b — ANSWERED, 2026-08-15.** The worry was that the OpenAI surface might drop
|
||||
reasoning the way the wiki documents for a mis-declared Anthropic backend. It does not. Checked at
|
||||
the API before any profile was switched:
|
||||
|
||||
| surface | request | result |
|
||||
|---|---|---|
|
||||
| `/anthropic/v1/messages` | `deepseek-v4-flash`, 64 tokens | 200, response carries a real `"type":"thinking"` block |
|
||||
| `/v1/chat/completions` | same | 200, message carries a populated `reasoning_content` (and a `reasoning` field) |
|
||||
| `/v1/models` | — | 200, exactly `["deepseek-v4-flash"]` — the exact-name trap is clear |
|
||||
| `/v1/models`, **no token** | — | **401** — Caddy is gating, as designed |
|
||||
|
||||
So reasoning survives on **both** surfaces, and the `/anthropic` choice for `local` is about protocol
|
||||
correctness rather than a repair for a known loss. The last row matters on its own: the wiki warns
|
||||
the gateway's own `SecurityPolicy` fails open, so it is worth knowing the proxy in front really does
|
||||
refuse an unauthenticated request here.
|
||||
|
||||
---
|
||||
|
||||
@@ -290,6 +304,160 @@ Merging config is not proving it. The checks, in order:
|
||||
|
||||
---
|
||||
|
||||
## 7.1 What the live run actually found — 2026-08-15
|
||||
|
||||
U1–U2c were done, the daemon restarted onto them, and both new profiles were spawned for real. The
|
||||
migration was then **reverted**. This section is the result, so none of it has to be re-derived.
|
||||
|
||||
### The blocker
|
||||
|
||||
`llm.ltms.dev` answers **HTTP 413 Request Entity Too Large** above **32 KiB (32768 bytes)**, on both
|
||||
surfaces:
|
||||
|
||||
```
|
||||
/v1 32695 bytes -> 200 /anthropic 32095 bytes -> 200
|
||||
/v1 32795 bytes -> 413 /anthropic 32855 bytes -> 413
|
||||
```
|
||||
|
||||
32 KiB is far below one real agent turn.
|
||||
|
||||
**Root cause — confirmed by the systems/vms side, 2026-08-15.** My guess that it was a Caddy
|
||||
`request_body max_size` was **wrong**. It is Envoy, inside `aigw` on `llm.vm`. Envoy Gateway defaults
|
||||
a listener's `per_connection_buffer_limit_bytes` to **32768**, and the AI Gateway buffers the *whole*
|
||||
request body before it can route on the model name — so that default is not a network tuning knob
|
||||
here, it is a hard ceiling on prompt size. Read out of the live Envoy `config_dump`:
|
||||
|
||||
```
|
||||
listener default/llm/http per_connection_buffer_limit_bytes: 32768
|
||||
```
|
||||
|
||||
Nobody chose 32 KiB; it was inherited from the default. Both TLS edges are innocent: the same
|
||||
boundary reproduces on the LAN path and the internet path, and both 413s carry an `x-llm-consumer`
|
||||
header their auth proxy sets only *after* authenticating — so the body cleared both edges and the
|
||||
auth. Directly on `llm.vm`, `aigw` 413s at 39 KB while the vLLM backend accepts the same 39 KB and
|
||||
answers 200.
|
||||
|
||||
**Do not plan around 32 KiB.** The intended ceiling is far higher. Their fix — a `ClientTrafficPolicy`
|
||||
setting `bufferLimit: 8Mi` — is written but **not deployed** as of this note, pending their operator's
|
||||
approval. I have not re-tested and will not until they confirm, so as not to measure a half-changed
|
||||
system. Fixed in **systems/vms**, not here.
|
||||
|
||||
### The part worth remembering
|
||||
|
||||
Two members were spawned at the same moment with the same message:
|
||||
|
||||
| | `local` (claude-code, `/anthropic`) | `gx` (opencode, `/v1`) |
|
||||
|---|---|---|
|
||||
| READY → BUSY | 19:07:26 | 19:07:45 |
|
||||
| BUSY → DONE | **19:08:51 (66s)** | **never — 10+ min, ticket FAILED** |
|
||||
|
||||
**`local` passed.** It passed only because the probe was three trivial questions in a fresh session,
|
||||
so the request fit under 32 KiB. The profile looked healthy and was a landmine set to fire on the
|
||||
first turn that reads a file.
|
||||
|
||||
So §7's checklist was not wrong, it was **too easy**. Any future run of it must use a task that reads
|
||||
a real file. A liveness probe proves the token and the URL; it does not prove the path.
|
||||
|
||||
`gx` did not fail loudly either. Reproduced outside the bridge by running `opencode` by hand with the
|
||||
launcher's own generated config:
|
||||
|
||||
```
|
||||
Error: Request Entity Too Large
|
||||
...compacts context, retries...
|
||||
Error: Request Entity Too Large
|
||||
```
|
||||
|
||||
opencode **catches the 413, compacts, and retries — indefinitely**. A member that fails loudly costs
|
||||
one turn; this one costs the whole task and is indistinguishable from a slow worker.
|
||||
|
||||
> **Diagnosing a stuck opencode member.** Do not read its pane. The launcher writes its config to a
|
||||
> temp dir and passes it as `OPENCODE_CONFIG` — find it with
|
||||
> `ls -dt /var/folders/*/*/T/bridged-opencode-* | head -1`, check the provider block and the key's
|
||||
> length and prefix (never its value), then reproduce with `opencode run --auto -m <provider>/<model>`
|
||||
> using the same `OPENCODE_CONFIG`. That is what turned "it hangs" into a one-line error.
|
||||
|
||||
### What checked out, and needs no re-testing
|
||||
|
||||
- Token accepted on both surfaces. **Unauthenticated → 401**, so the Caddy proxy really does gate —
|
||||
the wiki's "SecurityPolicy fails open" warning is about the gateway itself, not the edge.
|
||||
- `/v1/models` returns exactly `["deepseek-v4-flash"]`, so trap 3 is clear.
|
||||
- **Reasoning survives both surfaces** — see §3b above.
|
||||
- The launcher's generated opencode provider block is correct, carrying a real 48-character `llmk-`
|
||||
key rather than the `bridged-local-noauth` placeholder.
|
||||
- `SubscriptionGuard` accepted `llm.ltms.dev` after the allowlist edit and the restart: `local`
|
||||
spawned without throwing, which is the check that catches a missed restart.
|
||||
|
||||
### Resolution — both ceilings fixed, migration completed
|
||||
|
||||
systems/vms fixed both, and each was re-checked from this side rather than taken on trust:
|
||||
|
||||
| ceiling | was | now | our own check |
|
||||
|---|---|---|---|
|
||||
| listener buffer | 32 KiB | 32 Mi | 1.2 MB body → **200** (was 413) |
|
||||
| LLM route timeout | 60s | 86400s | the request that truncated: **101s, `message_stop` present, 4000/4000** |
|
||||
|
||||
The timeout moved in two steps on 2026-08-15: 60s → 1800s, then 1800s → **86400s (24 hours)** after
|
||||
the truncation risk below was discussed. They tried `request: 0s` first, which removes the
|
||||
total-duration timer completely. It works, but on an `AIGatewayRoute` the **idle timeout is derived
|
||||
from the request timeout**, so `0s` also removed any bound on a stalled connection. 86400s keeps a
|
||||
reaper for dead connections while putting the truncation timer out of practical reach.
|
||||
|
||||
Neither was deliberate. The 32 KiB was Envoy Gateway's default `per_connection_buffer_limit_bytes`;
|
||||
the 60s was Envoy AI Gateway's own documented default. The 60s bounded **generation** as well as
|
||||
prompt size — a tiny prompt with a long answer returned 504 at 60.05s.
|
||||
|
||||
Two configuration facts worth keeping, from their bisection:
|
||||
|
||||
- **`ClientTrafficPolicy` is honoured in standalone `aigw run`; `BackendTrafficPolicy` is NOT.** A
|
||||
`BackendTrafficPolicy` setting `requestTimeout` is accepted, logs nothing, and leaves the routes
|
||||
unchanged (upstream `envoyproxy/gateway#9513`). What works is `timeouts: {request: …}` on each
|
||||
`AIGatewayRoute` rule. Nothing from the outside distinguishes the two — the same silent-default
|
||||
shape as their `SecurityPolicy` caveat.
|
||||
- In that stack, "the config was accepted" proves nothing. Read the live `config_dump`.
|
||||
|
||||
## 7.2 The risk we accepted, and why we could not remove it
|
||||
|
||||
Raising the timeout made the failure **rare, not impossible**, and the residual failure is silent.
|
||||
|
||||
On a mid-response timeout over chunked HTTP/1.1, Envoy ends the chunked encoding *cleanly* instead of
|
||||
resetting the connection, so the client receives what looks like a complete transfer
|
||||
(`envoyproxy/envoy#17186` — acknowledged as a bug in 2021, closed by a stale bot, never fixed). The
|
||||
December 2025 fix `envoyproxy/envoy#42269` changes locally-originated resets from `NO_ERROR` to
|
||||
`INTERNAL_ERROR`, but it is **HTTP/2 only** and SSE clients here speak HTTP/1.1.
|
||||
|
||||
Measured on our side while the timeout was still 60s:
|
||||
|
||||
```
|
||||
HTTP 200 61.07s 141992 bytes
|
||||
message_stop 0 message_delta 0 error events 0
|
||||
emitted 2473 of 4000, ending on a WELL-FORMED SSE frame
|
||||
```
|
||||
|
||||
A syntactically valid stream that simply stops. Any timer firing mid-stream — route timeout, idle
|
||||
timeout, `max_stream_duration` — fails this same way.
|
||||
|
||||
**The recommended defence does not transfer to us.** The right fix is to treat a stream with no
|
||||
`message_stop` / `[DONE]` / `finish_reason` as failed. We cannot: our members are Claude Code and
|
||||
opencode, third-party clients whose SSE parsing we do not own, and there is no seam to insert the
|
||||
check. Whether either detects a missing terminator is unverified — and opencode's handling of the 413
|
||||
(swallow, compact, retry forever, never surface an error) does not suggest it is strict.
|
||||
|
||||
So the honest statement of our position:
|
||||
|
||||
> Gateway traffic is acceptable at 86400s because a single request would have to run for 24 hours to
|
||||
> trip the bug — **not** because we could detect it if it did.
|
||||
|
||||
At 86400s our **own** limit binds first, which is the ordering we want. `MessageService.ASYNC_TIMEOUT_MS`
|
||||
caps a turn at 30 minutes, so a runaway request ends as a clean `FAILED` ticket that we raised, rather
|
||||
than as a silently truncated `200` that we cannot see. While the gateway sat at 1800s the two numbers
|
||||
were equal and did not nest, so a gateway-side stall could have been misread as a bug in our own ticket
|
||||
handling. That ambiguity is now gone.
|
||||
|
||||
**If a member ever returns a confident but truncated answer, suspect this before anything in our own
|
||||
code.** That is the whole reason this section exists.
|
||||
|
||||
---
|
||||
|
||||
## 8. Related
|
||||
|
||||
- [CB-589 / #74](https://git.ltms.dev/lms/claude-bridge/issues/74) — cost-first placement and a
|
||||
|
||||
Reference in New Issue
Block a user