Compare commits

..

1 Commits

Author SHA1 Message Date
Dai Ha 0efe1567c0 CB-586: prune refs/wip/* older than 24h whose tree is reachable from main
CI / build (pull_request) Successful in 1m13s
CI / contract (pull_request) Successful in 1m21s
Add the CB-586 retention rule to GitWorktrees and drive it from the reaper:
a snapshot is deleted only when its tree content is already reachable from
main AND the ref is older than 24h. Reachability keeps the last copy of a
worker's work; the age floor stops a fresh snapshot being swept while a
lead is still looking at it. Every deletion logs the ref name and commit
sha so it is recoverable from the reflog. The /members response gains a
wipRefs{count,costBytes} census the operator can read without shelling
into the repo.
2026-08-16 19:02:33 +02:00
11 changed files with 746 additions and 395 deletions
@@ -8,8 +8,6 @@ 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;
@@ -18,41 +16,35 @@ import java.util.concurrent.TimeUnit;
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
* async delegation ticket ({@code bridge_send(wait:false)}) that reached a terminal phase
* (CB-588).
* 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.
*
* <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>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>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}.
* <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.
*/
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";
/** 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. */
/** CB-588: singular form, one uncollected ticket. */
static final String TICKET_NUDGE_FORMAT =
"Ticket %s finished%s — run bridge_poll(ticket=%s) to collect it";
/** Coalesced form, several uncollected tickets for the same lead. */
/** CB-588: 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";
@@ -64,11 +56,11 @@ 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<>();
/** Tickets that have gone terminal but not yet been polled, keyed by ticket. */
/** 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. */
private final ConcurrentHashMap<String, PendingTicket> pendingTickets = new ConcurrentHashMap<>();
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets). */
/** CB-588: leads with an active ticket-reminder schedule. */
private final ConcurrentHashMap<String, Boolean> activeLeads = new ConcurrentHashMap<>();
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
@@ -97,67 +89,178 @@ public final class ReplyPushLoop {
}
}
// --- pending-work lookups (package-private for unit-testing) -------------------------------
// --- decision logic (package-private for unit-testing) -------------------------------------
/** The action the loop should take for a lead at the given reminder count. */
/** The action the loop should take for a target at the given reminder count. */
enum Action { INJECT, WAIT_BUSY, STOP }
/**
* 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.
* 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
*/
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);
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;
}
return result;
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);
}
}
}
/** 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) {
}
/** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */
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());
/**
* 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);
}
/**
* 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
* 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.
*/
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);
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) {
return pendingTickets.values().stream().filter(t -> lead.equals(t.lead())).toList();
}
/**
* 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.
*/
Action decideTickets(String lead, int reminderCount) {
if (pendingFor(lead).isEmpty()) {
log.debug("push: nothing pending for lead {}, stopping ticket reminder", lead);
return Action.STOP;
}
if (reminderCount >= maxReminders) {
log.debug("push: reminder cap ({}) reached for lead {}, stopping", maxReminders, lead);
log.debug("push: ticket reminder cap ({}) reached for lead {}, stopping", maxReminders, lead);
countNudge("exhausted");
return Action.STOP;
}
@@ -175,185 +278,95 @@ public final class ReplyPushLoop {
return Action.WAIT_BUSY;
}
// --- public entrypoints ----------------------------------------------------------------------
/**
* 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());
}
/**
* 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>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);
/** 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 -> {
injectNudge(lead, reminderCount);
scheduleNext(lead, reminderCount + 1);
injectTicketNudge(lead, reminderCount);
scheduleTicketTick(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);
case WAIT_BUSY -> scheduleTicketTick(lead, reminderCount);
case STOP -> stopOrRestartTicketLoop(lead, pendingBefore);
}
}
/** 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());
}
/**
* 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).
* 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).
*
* <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>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>Package-private so a test can drive the interleaving directly rather than trying to force a
* genuine thread race.
* 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.
*
* <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.
* <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.
*/
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore) {
void stopOrRestartTicketLoop(String lead, Set<String> pendingBefore) {
activeLeads.remove(lead);
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);
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);
return;
}
log.debug("push: reminder loop ended for lead {}", lead);
log.debug("push: ticket reminder loop ended for lead {}", 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);
/** 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);
return;
}
String nudge = formatNudge(replyTargets, tickets);
String nudge = formatTicketsNudge(pending);
try {
agents.send(lead, nudge);
log.debug("push: nudge {}/{} sent to lead {} ({} reply target(s), {} ticket(s))",
reminderCount + 1, maxReminders, lead, replyTargets.size(), tickets.size());
log.debug("push: ticket nudge {}/{} sent to lead {} for {} ticket(s)",
reminderCount + 1, maxReminders, lead, pending.size());
countNudge("delivered");
} catch (RuntimeException e) {
log.warn("push: failed to nudge lead {} (reminder {}/{}): {}",
lead, reminderCount + 1, maxReminders, e.toString());
log.warn("push: failed to nudge lead {} for {} ticket(s) (reminder {}/{}): {}",
lead, pending.size(), reminderCount + 1, maxReminders, e.toString());
}
}
/** Schedule the next tick on the scheduler thread pool. */
private void scheduleNext(String lead, int nextReminderCount) {
scheduler.schedule(() -> tick(lead, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS);
/** 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);
}
// --- 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. */
/** Render one or several pending tickets as a single nudge line. */
private static String formatTicketsNudge(List<PendingTicket> pending) {
if (pending.size() == 1) {
PendingTicket t = pending.get(0);
@@ -370,22 +383,23 @@ public final class ReplyPushLoop {
// --- lifecycle -----------------------------------------------------------------------------
/**
* 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.
* 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.
*/
public boolean isActive() {
return !activeLeads.isEmpty();
return !activeTargets.isEmpty() || !activeLeads.isEmpty();
}
/** Shut down the scheduler. Outstanding reminders are cancelled. */
public void stop() {
scheduler.shutdownNow();
activeTargets.clear();
activeLeads.clear();
pendingReplies.clear();
pendingTickets.clear();
}
@@ -234,7 +234,14 @@ public final class BridgedApp {
List<Map<String, Object>> out = sessions.roster().stream()
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
.toList();
ctx.status(200).json(Map.of("workers", out));
Map<String, Object> body = new LinkedHashMap<>();
body.put("workers", out);
// CB-586: operator visibility for the refs/wip snapshot store without shelling into the
// repo — how many snapshot refs exist and roughly what they cost. Present only once a
// worktree session has established the repo, so a never-snapshotted fleet reports nothing.
sessions.wipRefs().ifPresent(st -> body.put("wipRefs",
Map.of("count", st.count(), "costBytes", st.costBytes())));
ctx.status(200).json(body);
}
/** The configured worker profiles and which one a no-argument spawn uses. */
@@ -12,9 +12,12 @@ import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.StandardCopyOption;
import java.security.SecureRandom;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.stream.Collectors;
@@ -299,6 +302,129 @@ public final class GitWorktrees implements Worktrees {
return index;
}
/**
* One {@code refs/wip/<branch>} snapshot ref as read by {@link #listWipRefs}: its full ref name,
* the snapshot commit's sha, and that commit's committer time in unix millis (the age of the
* snapshot — a snapshot is written once and never rewritten, so the commit date is the ref's).
*/
private record WipRef(String refName, String sha, long committerMillis) {
String branch() {
return refName.substring("refs/wip/".length());
}
}
@Override
public WipRefStats wipRefs(String repoRoot) {
List<WipRef> refs = listWipRefs(repoRoot);
long costBytes = 0;
for (WipRef ref : refs) {
costBytes += treeSize(repoRoot, ref.sha());
}
return new WipRefStats(refs.size(), costBytes);
}
@Override
public int pruneWipRefs(String repoRoot, long minAgeMillis) {
// The rule is documented on Worktrees#pruneWipRefs: delete only a snapshot whose tree
// content is already reachable from main AND that is older than minAgeMillis. Reachability
// is the floor that keeps a worker's last copy; the age floor keeps a just-written snapshot
// from being swept while a lead may still be looking at it.
List<WipRef> refs = listWipRefs(repoRoot);
if (refs.isEmpty()) {
return 0;
}
long nowMillis = System.currentTimeMillis();
// Resolve what main carries once per sweep, not once per ref.
Set<String> mainObjects = reachableObjectsFromMain(repoRoot);
int deleted = 0;
for (WipRef ref : refs) {
long ageMillis = nowMillis - ref.committerMillis();
if (ageMillis <= minAgeMillis) {
continue; // too recent — never swept, even if it looks recoverable (CB-586)
}
String tree = exec("git", "-C", repoRoot, "rev-parse", ref.sha() + "^{tree}").trim();
if (!mainObjects.contains(tree)) {
// Last copy of the snapshot's content — the worker's work exists nowhere else.
// Never delete automatically (CB-586 criterion 2).
continue;
}
exec("git", "-C", repoRoot, "update-ref", "-d", ref.refName());
deleted++;
log.info("pruned snapshot ref refs/wip/{} commit={} (age {}h): its tree is already "
+ "reachable from main, so the work is preserved; recover from reflog via "
+ "git update-ref refs/wip/{} {}",
ref.branch(), ref.sha(), TimeUnit.MILLISECONDS.toHours(ageMillis),
ref.branch(), ref.sha());
}
return deleted;
}
/**
* Every {@code refs/wip/*} ref (see {@link WipRef}). The committer date is read as a unix
* count of seconds and converted to millis. {@code %00} (NUL) separates the fields because a
* branch name may contain spaces.
*/
private List<WipRef> listWipRefs(String repoRoot) {
String out = exec("git", "-C", repoRoot, "for-each-ref",
"--format=%(refname)%00%(objectname)%00%(committerdate:unix)", "refs/wip/");
List<WipRef> refs = new ArrayList<>();
for (String line : out.split("\\R")) {
if (line.isBlank()) {
continue;
}
String[] parts = line.split("\u0000", -1);
if (parts.length == 3 && !parts[1].isBlank()) {
refs.add(new WipRef(parts[0], parts[1], Long.parseLong(parts[2]) * 1000L));
}
}
return refs;
}
/**
* The set of object shas reachable from {@code main}, or an empty set when {@code main} cannot
* be resolved. An empty set is the safe direction: the retention sweep then concludes nothing
* is recoverable, so it deletes nothing — a repo with no {@code main} must never cause a
* worker's last copy of a snapshot to be dropped on a reachability misreading.
*/
private Set<String> reachableObjectsFromMain(String repoRoot) {
if (exitCode("git", "-C", repoRoot, "rev-parse", "--verify", "main") != 0) {
log.debug("refs/wip retention: no 'main' ref in {} — treating nothing as reachable", repoRoot);
return Set.of();
}
String out = exec("git", "-C", repoRoot, "rev-list", "--objects", "main");
Set<String> objects = new HashSet<>();
for (String line : out.split("\\R")) {
if (line.isBlank()) {
continue;
}
int sp = line.indexOf(' ');
objects.add(sp < 0 ? line : line.substring(0, sp));
}
return objects;
}
/** Approximate cost of a snapshot: the sum of every blob's size in its committed tree. */
private long treeSize(String repoRoot, String sha) {
String out = exec("git", "-C", repoRoot, "ls-tree", "-r", "-l", sha);
long total = 0;
for (String line : out.split("\\R")) {
if (line.isBlank()) {
continue;
}
// ls-tree -l row: "<mode> <type> <object> <size>\t<path>"; the size is only numeric for
// blobs (trees read "-"), so gate on the type token and take the 4th whitespace field.
String[] parts = line.split("\\s+");
if (parts.length >= 4 && "blob".equals(parts[1])) {
try {
total += Long.parseLong(parts[3]);
} catch (NumberFormatException ignored) {
// a '-' size (or any anomaly) contributes nothing to the rough figure
}
}
}
return total;
}
/** Resolve the directory that will hold per-session worktree checkouts. */
private Path resolveRoot(String repoRoot) {
if (configuredRoot != null && !configuredRoot.isBlank()) {
@@ -53,6 +53,14 @@ public final class SessionManager implements TurnListener {
private final int contextCap;
private final boolean clearAfterTurn;
private volatile MemberLifecycle memberLifecycle = MemberLifecycle.NONE;
/**
* CB-586: the repo root the fleet actually works in, remembered the first time a worktree
* session is spawned (worktrees are checkouts of it). {@code refs/wip/*} live there, and this
* single cached value is what the snapshot retention sweep and the operator-visible census run
* against. The daemon is bridged into one project at a time, so "the first worktree's repo" is
* the repo; {@code null} until any worktree is spawned, meaning nothing to sweep or measure.
*/
private volatile String fleetRepoRoot;
/** CB-520: notified with a terminalId on every acquire; no-op until wired. */
private final List<Consumer<String>> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
@@ -450,6 +458,11 @@ public final class SessionManager implements TurnListener {
// The non-worktree path always used this chain; only this branch was missed.
String repoRoot = worktrees.repoRoot(
launcher.effectiveCwd(new SpawnRequest(preResolvedProfile, requestedCwd, callerCwd)));
if (fleetRepoRoot == null) {
// CB-586: remember the repo whose worktrees the fleet spawns — its refs/wip/* are the
// snapshot store the retention sweep and the operator census operate on.
fleetRepoRoot = repoRoot;
}
String branch = "worker/" + slug(wt.ticketSlug()) + "-" + nonce();
String path = null;
PeerHandle handle;
@@ -740,6 +753,26 @@ public final class SessionManager implements TurnListener {
return registry.size();
}
/**
* CB-586: the operator-visible census of {@code refs/wip/*} in the repo the fleet works in —
* how many snapshot refs exist and roughly what they cost. Empty (no repo known) until at
* least one worktree session has been spawned, exactly so a fleet that has never snapshotted
* anything surfaces nothing new, as it did before CB-586.
*/
public Optional<Worktrees.WipRefStats> wipRefs() {
String repo = fleetRepoRoot;
return repo == null ? Optional.empty() : Optional.of(worktrees.wipRefs(repo));
}
/**
* CB-586: run the snapshot retention sweep in the fleet's repo (a no-op until a worktree has
* been spawned, which establishes the repo). Returns how many {@code refs/wip/*} it deleted.
*/
public int sweepWipRefs(long minAgeMillis) {
String repo = fleetRepoRoot;
return repo == null ? 0 : worktrees.pruneWipRefs(repo, minAgeMillis);
}
/**
* The registered session owning {@code terminalId}, or {@code null} if none does.
*
@@ -14,12 +14,20 @@ public final class SessionReaper {
private static final Logger log = LoggerFactory.getLogger(SessionReaper.class);
private static final long DEFAULT_INTERVAL_MILLIS = 5000;
/** CB-586: the refs/wip age floor — never sweep a snapshot younger than 24h (the CB-586 rule). */
private static final long WIP_MIN_AGE_MILLIS = TimeUnit.HOURS.toMillis(24);
/**
* CB-586: how often the retention sweep runs. Given the 24h age floor, running it every few
* hours means a ref is dropped within hours of becoming eligible, never within minutes.
*/
private static final long WIP_SWEEP_INTERVAL_NANOS = TimeUnit.HOURS.toNanos(6);
private final SessionManager sessions;
private final long idleTtlNanos;
private final long intervalMillis;
private volatile boolean running;
private Thread thread;
private volatile long lastWipSweepNanos = Long.MIN_VALUE;
/** Construct a reaper with the default 5-second polling interval. */
public SessionReaper(SessionManager sessions, long idleTtlSeconds) {
@@ -49,10 +57,32 @@ public final class SessionReaper {
} catch (RuntimeException e) {
log.warn("session reaper iteration failed; continuing", e);
}
maybeSweepWipRefs();
sleep();
}
}
/**
* CB-586: run the refs/wip retention sweep on a slow cadence (hours, not the per-iteration
* millisecond loop). Best-effort — a failure must never take the idle-reap loop down with it.
*/
private void maybeSweepWipRefs() {
long now = System.nanoTime();
if (now - lastWipSweepNanos < WIP_SWEEP_INTERVAL_NANOS) {
return;
}
try {
int deleted = sessions.sweepWipRefs(WIP_MIN_AGE_MILLIS);
if (deleted > 0) {
log.info("refs/wip retention sweep deleted {} snapshot ref(s) older than 24h whose "
+ "content was already reachable from main", deleted);
}
} catch (RuntimeException e) {
log.warn("refs/wip retention sweep failed; continuing", e);
}
lastWipSweepNanos = now;
}
private void sleep() {
try {
Thread.sleep(intervalMillis);
@@ -54,4 +54,48 @@ public interface Worktrees {
* tolerance — a worktree that is gone holds nothing to snapshot)
*/
Optional<String> snapshot(String worktreePath, String branch, String message);
/**
* CB-586: how many {@code refs/wip/*} snapshot refs exist in {@code repoRoot} and roughly what
* they cost. This is the operator-visible surface for the snapshot growth CB-578 stage C left
* behind — counts of refs alone hide that each one pins a whole tree for {@code git gc}.
*
* @param repoRoot the repository to scan
* @return count of snapshot refs, and {@code costBytes} = the approximate total working-tree
* size of every snapshot's committed content (summed per ref, so shared objects are
* counted once per ref that carries them)
*/
WipRefStats wipRefs(String repoRoot);
/**
* CB-586: run the {@code refs/wip/*} retention sweep and return how many refs it deleted.
*
* <p>The retention rule is <em>reachability plus an age floor</em>. A snapshot ref is deleted
* only when <strong>both</strong> hold:
* <ol>
* <li>its commit's <em>tree content</em> is already reachable from {@code main} — the work
* the snapshot preserved has been recovered, so dropping the ref loses nothing; and</li>
* <li>the ref is older than {@code minAgeMillis} — a very recent snapshot is never swept
* while a lead may still be looking at it.</li>
* </ol>
*
* <p>Reachability is the safety property. A snapshot exists precisely because the work was not
* committed anywhere else, so a snapshot whose content is <em>not</em> reachable from
* {@code main} is the <strong>last copy</strong> of a worker's work and must never be deleted
* automatically — that is the failure CB-576 and CB-578 stage C were built to stop. Age alone
* must never drive a deletion, because age-based sweeping is exactly how the last copy gets
* destroyed. (Both numbers and the rule are CB-586's decision; this method only implements it.)
*
* <p>Every deletion logs the ref name and the commit sha, so an operator who finds they lost
* the wrong thing can still recover it from git's reflog.
*
* @param repoRoot the repository whose {@code refs/wip/*} to sweep
* @param minAgeMillis the age floor; a ref younger than this is never touched
* @return the number of snapshot refs deleted
*/
int pruneWipRefs(String repoRoot, long minAgeMillis);
/** CB-586: the operator-visible census of {@code refs/wip/*} in one repository. */
record WipRefStats(int count, long costBytes) {
}
}
@@ -20,7 +20,6 @@ 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.*;
@@ -28,12 +27,6 @@ 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.
@@ -65,50 +58,38 @@ class ReplyPushLoopTest {
// --- decide() logic ------------------------------------------------------------------------
@Test
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");
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));
}
@Test
void decideWithNothingPendingIsStop() {
void decideWithEmptyInboxIsStop() {
agents = agentWithStatus("idle");
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
}
@Test
void decideAtCapIsStop() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
var loop = loop(2, 100);
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100).decide(WORKER, 2));
}
@Test
void decideUnderCapWithInjectablePrimaryIsInject() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
}
@Test
void decideUnderCapWithBlockedPrimaryIsInject() {
agents = agentWithStatus("blocked");
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0),
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
"BLOCKED is injectable");
}
@@ -116,9 +97,7 @@ class ReplyPushLoopTest {
void decideUnderCapWithDonePrimaryIsInject() {
agents = agentWithStatus("done");
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0),
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0),
"DONE is injectable");
}
@@ -126,29 +105,23 @@ class ReplyPushLoopTest {
void decideUnderCapWithBusyPrimaryIsWaitBusy() {
agents = agentWithStatus("working");
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
}
@Test
void decideUnderCapWithUnknownPrimaryIsWaitBusy() {
agents = agentWithStatus("unknown");
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop().decide(WORKER, 0));
}
@Test
void decideStopsAfterInboxIsEmptied() {
agents = agentWithStatus("idle");
inbox.publish(WORKER, "m1", "hello");
var loop = loop();
loop.onReplyQueued(WORKER);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.INJECT, loop().decide(WORKER, 0));
inbox.ack(WORKER, "m1");
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(WORKER, 0));
}
// --- onReplyQueued integration -------------------------------------------------------------
@@ -228,19 +201,12 @@ class ReplyPushLoopTest {
assertTrue(nudge.contains("bridge_poll(target=term_worker)"));
}
@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 -----------------------
// --- CB-588: async ticket terminal nudges — decideTickets() logic --------------------------
@Test
void decideTicketsWithNothingPendingIsStop() {
agents = agentWithStatus("idle");
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop().decideTickets(PRIMARY, 0));
}
@Test
@@ -248,7 +214,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
var loop = loop(2, 100_000);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2));
}
@Test
@@ -256,7 +222,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("idle");
var loop = loop(5, 100_000);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.INJECT, loop.decideTickets(PRIMARY, 0));
}
@Test
@@ -264,7 +230,7 @@ class ReplyPushLoopTest {
agents = agentWithStatus("working");
var loop = loop(5, 100_000);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decideTickets(PRIMARY, 0));
}
@Test
@@ -272,7 +238,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.decide(PRIMARY, 0));
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 0));
}
// --- CB-588: onTicketTerminal integration ---------------------------------------------------
@@ -289,7 +255,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-only nudge must not tell the lead to run the target-poll call");
assertFalse(nudge.contains("bridge_poll(target="), "a ticket nudge must not tell the lead to run the target-poll call");
}
@Test
@@ -380,11 +346,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
// decide's STOP releases it, so the ticket coalesces onto a schedule that is about to
// decideTickets' 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 stopOrRestart — the STOP path's own release-and-recheck —
// reliable, so this drives stopOrRestartTicketLoop — 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 (ticketsBefore=empty), standing in for one
// that was NOT part of the pre-decision snapshot (pendingBefore=empty), standing in for one
// that races in during the decision-to-release window.
var rec = recordingClient();
agents = new AgentControl(rec);
@@ -392,9 +358,9 @@ class ReplyPushLoopTest {
loop.onTicketTerminal("task-1", WORKER, false); // pendingTickets={task-1}; activeLeads={PRIMARY}
// Stand in for the scheduler thread reaching decide==STOP for this lead — with nothing
// Stand in for the scheduler thread reaching decideTickets==STOP for this lead — with nothing
// pending at decide time — while task-1 races in before the release below runs.
loop.stopOrRestart(PRIMARY, Set.of(), Set.of());
loop.stopOrRestartTicketLoop(PRIMARY, 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");
@@ -402,19 +368,18 @@ class ReplyPushLoopTest {
@Test
void aStaleUncollectedTicketAtCapDoesNotRestartTheLoop() {
// The other direction of the same fix: restarting on ANY non-empty pending set would be
// The other direction of the same fix: restarting on ANY non-empty pendingFor(lead) 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
// #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.
// #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.
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.stopOrRestart(PRIMARY, Set.of(), Set.of("task-1"));
loop.stopOrRestartTicketLoop(PRIMARY, 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");
@@ -437,52 +402,7 @@ class ReplyPushLoopTest {
void inboxNudgeStillUsesTheOriginalTargetPollCall() {
String nudge = ReplyPushLoop.NUDGE_FORMAT.formatted(WORKER, WORKER);
assertTrue(nudge.contains("bridge_poll(target=" + WORKER + ")"),
"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);
"CB-588 must not change the CB-307 inbox nudge's call shape");
}
// --- metrics (CB-512) ----------------------------------------------------------------------
@@ -510,10 +430,8 @@ 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.decide(PRIMARY, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop(2, 100, metrics).decide(WORKER, 2));
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the reminder cap must count as exhausted");
@@ -541,7 +459,7 @@ class ReplyPushLoopTest {
var loop = loop(2, 100_000, metrics);
loop.onTicketTerminal("task-1", WORKER, false);
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
assertEquals(ReplyPushLoop.Action.STOP, loop.decideTickets(PRIMARY, 2));
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
"hitting the ticket reminder cap must count as exhausted");
@@ -629,49 +547,4 @@ 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() {
}
}
}
@@ -27,11 +27,15 @@ public final class FakeWorktrees implements Worktrees {
public record SnapshotCall(String worktreePath, String branch, String message) {
}
public record PruneCall(String repoRoot, long minAgeMillis) {
}
private final List<AddCall> addCalls = new CopyOnWriteArrayList<>();
private final List<RemoveCall> removeCalls = new CopyOnWriteArrayList<>();
private final List<OverlayCall> overlayCalls = new CopyOnWriteArrayList<>();
private final List<RepoRootCall> repoRootCalls = new CopyOnWriteArrayList<>();
private final List<SnapshotCall> snapshotCalls = new CopyOnWriteArrayList<>();
private final List<PruneCall> pruneCalls = new CopyOnWriteArrayList<>();
private final Set<String> existingPaths = ConcurrentHashMap.newKeySet();
private final Set<String> trackedPaths = ConcurrentHashMap.newKeySet();
private final AtomicLong snapshotSeq = new AtomicLong();
@@ -40,6 +44,8 @@ public final class FakeWorktrees implements Worktrees {
private volatile boolean dirty = false;
private volatile String repoRoot = "/repo";
private volatile String prefix = "/worktrees";
private volatile WipRefStats wipRefs = new WipRefStats(0, 0L);
private volatile int pruneResult = 0;
public FakeWorktrees withRepoRoot(String root) {
this.repoRoot = root;
@@ -82,6 +88,18 @@ public final class FakeWorktrees implements Worktrees {
return this;
}
/** Configure the value returned by {@link #wipRefs}. */
public FakeWorktrees withWipRefs(WipRefStats stats) {
this.wipRefs = stats;
return this;
}
/** Configure the value returned by {@link #pruneWipRefs}. */
public FakeWorktrees withPruneResult(int deleted) {
this.pruneResult = deleted;
return this;
}
@Override
public String add(String repoRoot, String branch, String baseRef) {
addCalls.add(new AddCall(repoRoot, branch, baseRef));
@@ -135,6 +153,17 @@ public final class FakeWorktrees implements Worktrees {
return Optional.of("wip" + snapshotSeq.incrementAndGet());
}
@Override
public WipRefStats wipRefs(String repoRoot) {
return wipRefs;
}
@Override
public int pruneWipRefs(String repoRoot, long minAgeMillis) {
pruneCalls.add(new PruneCall(repoRoot, minAgeMillis));
return pruneResult;
}
public List<AddCall> addCalls() {
return List.copyOf(addCalls);
}
@@ -167,6 +196,10 @@ public final class FakeWorktrees implements Worktrees {
return List.copyOf(snapshotCalls);
}
public List<PruneCall> pruneCalls() {
return List.copyOf(pruneCalls);
}
public SnapshotCall lastSnapshot() {
return snapshotCalls.isEmpty() ? null : snapshotCalls.getLast();
}
@@ -3,6 +3,7 @@ package dev.ltms.bridged.session;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.HashSet;
@@ -138,6 +139,53 @@ class GitWorktreesTest {
return out;
}
/** Write {@code content} as a blob into the object database; returns its sha. */
private static String blobOf(Path cwd, String content) throws Exception {
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "hash-object", "-w", "--stdin")
.redirectErrorStream(true).start();
p.getOutputStream().write(content.getBytes(StandardCharsets.UTF_8));
p.getOutputStream().close();
String out = new String(p.getInputStream().readAllBytes()).trim();
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git hash-object timed out");
assertEquals(0, p.exitValue(), "git hash-object failed:\n" + out);
return out;
}
/** Build a single-file tree object from {@code blob}; returns the tree's sha. */
private static String treeOf(Path cwd, String path, String blob) throws Exception {
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "mktree")
.redirectErrorStream(true).start();
p.getOutputStream().write(("100644 blob " + blob + "\t" + path + "\n").getBytes(StandardCharsets.UTF_8));
p.getOutputStream().close();
String out = new String(p.getInputStream().readAllBytes()).trim();
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git mktree timed out");
assertEquals(0, p.exitValue(), "git mktree failed:\n" + out);
return out;
}
/** {@code git commit-tree} rooted at {@code tree} with a chosen committer date; returns the sha. */
private static String commitTree(Path cwd, String tree, String parent, String committerDate,
String message) throws Exception {
ProcessBuilder pb = new ProcessBuilder("git", "-C", cwd.toString(), "commit-tree",
tree, "-p", parent, "-m", message);
pb.environment().put("GIT_COMMITTER_DATE", committerDate);
Process p = pb.redirectErrorStream(true).start();
String out = new String(p.getInputStream().readAllBytes()).trim();
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git commit-tree timed out");
assertEquals(0, p.exitValue(), "git commit-tree failed:\n" + out);
return out;
}
/** {@code git update-ref <ref> <sha>} — create the snapshot ref directly. */
private static void updateRef(Path cwd, String ref, String sha) throws Exception {
git(cwd, "update-ref", ref, sha);
}
/** True when {@code ref} exists in the repo (for-each-ref on a missing ref is empty, not an error). */
private static boolean refExists(Path cwd, String ref) throws Exception {
return !forEachRef(cwd, ref).trim().isEmpty();
}
/**
* The heart of CB-525: a provisioned worktree must not inherit the primary's MCP servers. Without
* the isolation step the checked-out {@code .mcp.json} carries them in, and a worker navigating
@@ -474,4 +522,104 @@ class GitWorktreesTest {
assertThrows(WorktreeException.class, () -> gitWorktrees.snapshot(wt, branch, "test snapshot"),
"an unresolvable real index must fail loudly, not silently snapshot from an empty index");
}
/**
* CB-586, criterion 2. A snapshot whose content is NOT reachable from {@code main} is the last
* copy of a worker's work, and must never be deleted automatically — even when it is old and
* even when the caller passes a zero age floor. Uses the real snapshot path on a dirty worktree,
* so the unreachable tree is exactly the shape CB-576/CB-578 stage C exist to protect.
*/
@Test
void anUnreachableSnapshotIsNeverPruned(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
String branch = "cb-586-unreachable";
String wt = gitWorktrees.add(repo.toString(), branch, "HEAD");
Files.writeString(Path.of(wt).resolve("worker-draft.txt"), "work that exists nowhere else\n");
Optional<String> ref = gitWorktrees.snapshot(wt, branch, "snapshot with unreachable content");
assertTrue(ref.isPresent());
// Age floor 0 makes age a non-issue: only reachability can save it — and it must.
assertEquals(0, gitWorktrees.pruneWipRefs(repo.toString(), 0),
"the unreachable snapshot is the last copy and must not be pruned");
assertTrue(refExists(repo, "refs/wip/" + branch),
"an unreachable snapshot must survive the sweep");
}
/**
* CB-586, criterion 1 (the reachable half). A snapshot whose tree content IS already reachable
* from {@code main} and which is older than the age floor is pure duplication — the work is
* recovered — so it must be pruned.
*/
@Test
void aReachableSnapshotOlderThanTheFloorIsPruned(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
// A snapshot whose tree is exactly main's current tree: fully reachable from main.
String mainTree = revParse(repo, "main^{tree}");
String old = commitTree(repo, mainTree, revParse(repo, "HEAD"), "2020-01-01T00:00:00", "snapshot");
updateRef(repo, "refs/wip/recovered", old);
assertEquals(1, gitWorktrees.pruneWipRefs(repo.toString(), TimeUnit.HOURS.toMillis(24)),
"an old, main-reachable snapshot must be pruned");
assertFalse(refExists(repo, "refs/wip/recovered"),
"the reachable snapshot's ref must be gone after the sweep");
}
/**
* CB-586, the age floor. A snapshot whose content IS reachable from {@code main} but which is
* younger than the age floor must not be swept — a lead may still be looking at it.
*/
@Test
void aReachableButRecentSnapshotIsNotPruned(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
// Reachable from main, but committed "now" — a fresh snapshot. The 24h floor must protect it.
String mainTree = revParse(repo, "main^{tree}");
String fresh = commitTree(repo, mainTree, revParse(repo, "HEAD"),
"2038-01-01T00:00:00", "snapshot just taken");
updateRef(repo, "refs/wip/fresh", fresh);
assertEquals(0, gitWorktrees.pruneWipRefs(repo.toString(), TimeUnit.HOURS.toMillis(24)),
"a recent snapshot must be kept even when reachable");
assertTrue(refExists(repo, "refs/wip/fresh"),
"the recent reachable snapshot must survive the sweep");
}
/**
* CB-586, criterion 5. A fleet that has never snapshotted anything has no {@code refs/wip/*},
* so a sweep is a no-op and the census reports none — identical to before CB-586 existed.
*/
@Test
void aFleetWithNoSnapshotsPrunesNothingAndReportsNothing(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
assertEquals(0, gitWorktrees.pruneWipRefs(repo.toString(), 0),
"no snapshot refs means nothing to prune");
Worktrees.WipRefStats stats = gitWorktrees.wipRefs(repo.toString());
assertEquals(0, stats.count(), "a never-snapshotted fleet has zero refs/wip refs");
assertEquals(0L, stats.costBytes(), "a never-snapshotted fleet costs zero bytes");
}
/** CB-586, criterion 4: the census reports how many refs exist and roughly what they cost. */
@Test
void wipRefsReportsCountAndCost(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
String blob = blobOf(repo, "a recoverable snapshot's worth of content");
String tree = treeOf(repo, "snapshot.txt", blob);
updateRef(repo, "refs/wip/one", commitTree(repo, tree, revParse(repo, "HEAD"),
"2020-01-01T00:00:00", "snapshot"));
updateRef(repo, "refs/wip/two", commitTree(repo, tree, revParse(repo, "HEAD"),
"2020-01-02T00:00:00", "snapshot"));
Worktrees.WipRefStats stats = gitWorktrees.wipRefs(repo.toString());
assertEquals(2, stats.count(), "two snapshot refs are reported");
assertTrue(stats.costBytes() > 0, "the cost of the snapshots is a positive byte count");
}
}
@@ -138,6 +138,16 @@ class SessionManagerTest {
return java.util.Optional.of("wip" + snapshotSeq.incrementAndGet());
}
@Override
public WipRefStats wipRefs(String repoRoot) {
return new WipRefStats(0, 0L);
}
@Override
public int pruneWipRefs(String repoRoot, long minAgeMillis) {
return 0;
}
List<String> removeCalls() {
return List.copyOf(removeCalls);
}
@@ -444,4 +444,37 @@ class WorktreeSessionManagerTest {
"a failed snapshot leaves no ref to report");
}
/**
* CB-586, criteria 4 and 1. Once a worktree session is spawned the repo is known, so the
* operator census and the retention sweep delegate to that repo's {@code refs/wip/*}.
*/
@Test
void wipRefsAndSweepDelegateToTheFleetRepoOnceKnown() {
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
.withWipRefs(new Worktrees.WipRefStats(3, 42L)).withPruneResult(2);
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
assertTrue(sessions.wipRefs().isEmpty(),
"no worktree spawned yet means no repo is known and nothing to report");
assertEquals(0, sessions.sweepWipRefs(TimeUnit.HOURS.toMillis(24)),
"no worktree spawned yet means the sweep is a no-op");
sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-586", null));
Worktrees.WipRefStats stats = sessions.wipRefs().orElseThrow();
assertEquals(3, stats.count(), "the census comes from the fleet repo");
assertEquals(42L, stats.costBytes(), "the cost comes from the fleet repo");
assertEquals(2, sessions.sweepWipRefs(TimeUnit.HOURS.toMillis(24)),
"the sweep runs against the fleet repo");
List<FakeWorktrees.PruneCall> prunes = worktrees.pruneCalls();
assertEquals(1, prunes.size(), "the no-op short-circuits before reaching the seam, so only "
+ "the repo-known sweep issues a call");
assertEquals("/repo", prunes.getFirst().repoRoot(), "the sweep targets the fleet repo");
assertEquals(TimeUnit.HOURS.toMillis(24), prunes.getFirst().minAgeMillis(),
"the caller's age floor is passed through");
}
}