Compare commits

...

44 Commits

Author SHA1 Message Date
Dai Ha c393600921 CB-576: prove teardown survives an already-gone worktree 2026-08-15 09:31:10 +02:00
Dai Ha b525b0f08f CB-576: hasUncommitted tolerates an already-gone worktree
CI / contract (pull_request) Successful in 1m8s
CI / build (pull_request) Successful in 1m37s
2026-08-15 08:47:12 +02:00
Dai Ha 9118ce2537 CB-576: release preserves a dirty worktree instead of deleting it
CI / contract (pull_request) Successful in 1m2s
CI / build (pull_request) Successful in 1m35s
2026-08-15 07:34:06 +02:00
Dai Ha 2f48e08f1f Merge CB-577 follow-up: drop the target-keyed async index
CI / contract (push) Successful in 43s
CI / build (push) Successful in 55s
asyncTasksByWaiter correlates an async question by the exact rendezvous
waiter, so the target-keyed set it replaced can no longer decide
anything. Keeping it meant two indexes of the same fact, one of them
ambiguous whenever a target has two accepted tickets.

The race the old test modelled by reflection is gone with it: identity
keys make 'some other task reached this target' unrepresentable, so
there is no longer a wrong task for the question to land on.

Also carries the criterion-1 doc correction, which is identical to
aac29d6 on a different parent.
2026-08-15 06:40:31 +02:00
Dai Ha fec284e7cb Merge M4 unit 2a: accepted-turn identity carried to delivery
CI / build (push) Successful in 59s
CI / contract (push) Successful in 1m17s
TurnToken, owned by MessageService, binds a target to the exact
rendezvous waiter for one accepted send. Injector.Pending carries it and
the delivery callback hands it to CompletionResolver, so the baseline is
bound to the send it belongs to by construction rather than by a lookup
that could pick a different one.

The callback signature is required, not a defaulted overload: a delivery
with no token is exactly the unbound baseline this unit forbids, so a
default would let a caller silently produce it.

The token deliberately omits the session turn number. MessageService
owns acceptance but never learns of delivery, and
CompletionResolver.onDelivered runs before SessionManager.onDelivered,
so the number does not exist yet at the only point the token could
capture it. docs/M4-Fleet-Health.md criterion 1 records this and the two
rejected alternatives.

Still open for the next slice: the missing/post-restart baseline test and
the no-replay test.
2026-08-15 06:36:46 +02:00
Dai Ha 8e2e4c5e73 M4 unit 2a: migrate test call sites to the required turn token
The delivery callback now requires a TurnToken, so 54 test call sites
had to pass one. They use an explicit TestTurnTokens.inert(target)
rather than a defaulted overload, because a delivery with no token is
the unbound baseline this unit forbids.

The first version of inert() returned a fresh CompletableFuture as the
waiter, which turned captureBaselineSkipsTheReadWhenNoSendIsWaiting red:
the resolver saw a non-null waiter, concluded a turn was in flight, and
scraped a pane no send was blocked on. An inert value must omit the
fact, not invent it, so the waiter is now null and the production skip
fires as designed.
2026-08-15 06:36:35 +02:00
Dai Ha 0edc6615fc M4: correct unit 2 criterion 1 — TurnToken cannot carry the session turn
The criterion required the token to bind the session turn number. Three
independent refusals from the implementer showed why that is not
implementable at this layer: MessageService owns acceptance but never
learns of delivery, and CompletionResolver.onDelivered runs before
SessionManager.onDelivered, so the turn number does not exist yet at the
only point the token could capture it.

Records both rejected alternatives and why, so the next reader does not
re-derive them: a target-keyed registry restores the ambiguity the token
exists to remove, and injecting a turn counter couples layers to fill a
field nothing reads yet.
2026-08-15 06:32:02 +02:00
Dai Ha b745e159de CB-573: carry accepted turn tokens on delivery 2026-08-15 06:30:41 +02:00
Dai Ha aac29d604c M4: correct unit 2 criterion 1 — TurnToken cannot carry the session turn
CI / contract (push) Successful in 1m0s
CI / build (push) Successful in 1m32s
The criterion required the token to bind the session turn number. Three
independent refusals from the implementer showed why that is not
implementable at this layer: MessageService owns acceptance but never
learns of delivery, and CompletionResolver.onDelivered runs before
SessionManager.onDelivered, so the turn number does not exist yet at the
only point the token could capture it.

Records both rejected alternatives and why, so the next reader does not
re-derive them: a target-keyed registry restores the ambiguity the token
exists to remove, and injecting a turn counter couples layers to fill a
field nothing reads yet.
2026-08-15 06:29:27 +02:00
Dai Ha c884802b13 CB-577: remove obsolete async target tracking
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 56s
2026-08-15 06:28:53 +02:00
Dai Ha 5f5573a24e Merge CB-577: correlate an async question by its exact waiter
CI / contract (push) Failing after 0s
CI / build (push) Successful in 1m13s
markAsyncQuestion picked the first not-done task out of an unordered
set, so between resolveQuestion waking the first async send and the
question being recorded, a queued second send could join the set and
take the question. A lead answering with bridge_send{turnId} would then
resume a turn it did not mean to.

Each async task is now indexed by its exact rendezvous waiter, which has
identity semantics, so no other task can hold the same key. The question
is recorded before resolveQuestion, with a rollback when no waiter is
there, which closes the window rather than narrowing it.

An unanswered async question stays PENDING — the worker resumes after
its ask times out, so the delegation is not failed — and its stale
target tracking is now cleared instead of leaking.
2026-08-15 06:27:16 +02:00
Dai Ha 74b0087ebb CB-577: model async question ownership race
CI / build (pull_request) Successful in 1m29s
CI / contract (pull_request) Failing after 0s
2026-08-15 06:24:05 +02:00
Dai Ha 927e0151d4 CB-577: test async question waiter ownership 2026-08-15 06:23:22 +02:00
Dai Ha 5275922d1d CB-577: handle questions without async waiters 2026-08-15 06:22:22 +02:00
Dai Ha e186c7945a CB-577: correlate async questions to turns 2026-08-15 06:21:53 +02:00
Dai Ha 33a6e77f0e Merge CB-573: dormant fleet health monitor (M4 unit 1)
CI / build (push) Successful in 49s
CI / contract (push) Successful in 1m14s
An opt-in whole-fleet observer, separate from the 250ms delivery poller.
One AgentControl.list and one roster snapshot per tick, joined and fed to
the FleetHealth classifier, because a fault is a disagreement between the
two views at the same instant. Absent a health: block nothing is built
and no herdr call is made.

Adds bridge_list healthCoverage: off, detection-only, or full. Detection
is deliberately separate from notification, so a single-lead setup with
no webhook still gets detection and is told its coverage is partial
rather than being refused.

Two review fixes worth naming. tick() rescheduled itself as its last
statement with no try/catch, and a ScheduledExecutorService does not
re-run a task that threw — so the first agents.list failure would have
stopped health permanently and silently, which is exactly when the
control link is down. It now catches Throwable and reschedules in a
finally. And the snapshot fields this unit cannot supply are the named
constant NOT_YET_OBSERVED rather than bare false literals, because false
means no fault to this classifier.
2026-08-15 06:19:15 +02:00
Dai Ha 4e47489d53 Merge CB-568c: fail every pending async ticket on teardown
A released target left its second async ticket pending for the full
30-minute async timeout. abandon resolved only the rendezvous waiter,
and async tickets live in a separate map that could not even represent
two tasks on one target.

abandon now sweeps every non-question async ticket for the target, and
asyncTasksByTarget holds a set. The sweep is a plain loop: the first
attempt used Stream.anyMatch, which short-circuits on the first true, so
it completed one ticket and left the rest pending — the exact bug it was
fixing. Its test passed only because the rendezvous path failed the
first ticket anyway; the test now uses three tickets so a single
completion cannot satisfy it.

An ASKING ticket is an active turn, not a pending send, so the sweep
skips it and CB-574 is unaffected.
2026-08-15 06:18:28 +02:00
Dai Ha c76b2f149e CB-573: keep health monitoring after failures
CI / contract (pull_request) Successful in 42s
CI / build (pull_request) Successful in 1m15s
2026-08-15 06:17:28 +02:00
Dai Ha 826fffe05b CB-573: add dormant fleet health monitor 2026-08-15 06:15:24 +02:00
Dai Ha 75f57cdba7 CB-568c: fail every queued async ticket
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 1m34s
2026-08-15 06:15:21 +02:00
Dai Ha f556af5d4e CB-568: fail queued async tickets on teardown 2026-08-15 06:14:28 +02:00
Dai Ha d5f33f0c6e Merge CB-568: a dropped send reports the real cause
CI / build (push) Successful in 55s
CI / contract (push) Successful in 1m0s
Injector.drop knew the precise cause (herdr agent_not_found) but the
sender was told only 'worker unreachable or stuck', so a lead could not
tell a dead pane from a stalled model.

TurnListener.onTurnFailed gains a reason, defaulting to the old one-arg
form. CompletionResolver prefers that reason, then the pane scrape, then
the old fixed text.

drop now fires onTurnFailed unconditionally. That is the substantive
fix: the sender blocks on the rendezvous waiter, not on the delivered
future, so failing delivered() alone never woke it and a queued send sat
until its timeout.
2026-08-15 06:08:41 +02:00
Dai Ha 16e17b32ad CB-568: preserve dropped turn causes
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 1m26s
2026-08-15 06:06:07 +02:00
Dai Ha 988494e18b Merge M4 fleet-health design (docs only)
The full 5-unit design behind M4: evidence model, classification
precedence, the automatic-vs-lead action boundary, worktree safety on
release, typed inbox and lead routing, capacity, and human escalation.

Two decisions worth keeping visible. Detection is split from
notification, so health works in a single-lead setup with no webhook and
reports partial coverage instead of refusing to run. And capacity stays
a view: the bridge reports free slots but never spawns, reassigns, or
stops a member to improve utilisation, because only the lead holds the
work list.

Section 13 records eleven things nobody checked, including live LavinMQ,
OpenCode pane fixtures, and multi-lead routing.
2026-08-15 06:05:26 +02:00
Dai Ha c4549a5e20 Merge CB-575: filter the routine MCP cancellation WARN
The MCP SDK 2.0.0 registers no handler for notifications/cancelled, so every
client abort logged a WARN. M4 fleet health treats WARN as action-needed, so
that noise had a cost. A Logback TurboFilter denies only that one event:
right logger, WARN level, the SDK's exact format string, and a
JSONRPCNotification whose method is notifications/cancelled. Everything else
is NEUTRAL. If a later SDK handles cancellation the filter stops matching.
2026-08-15 06:00:40 +02:00
Dai Ha 46fa4f38d5 CB-575: filter MCP cancellation warnings
CI / build (pull_request) Successful in 49s
CI / contract (pull_request) Successful in 1m1s
2026-08-15 05:57:45 +02:00
Dai Ha 20e0e68ad7 M4: align design status with CB-573
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 1m28s
2026-08-15 05:56:00 +02:00
Dai Ha 17468a234a M4: document fleet health design 2026-08-15 05:55:00 +02:00
Dai Ha 8a53d5bfc6 Merge CB-574: an async delegation can now receive a worker's question
CI / contract (push) Successful in 1m2s
CI / build (push) Failing after 1m40s
A worker on a wait:false delegation called bridge_ask and the lead never
saw the question. Outcome.QUESTION is deliberately non-terminal, but
taskView tested r.completed() and fell into the failure branch, so the
ticket was marked FAILED and both the question text and its turnId were
discarded. The worker blocked for 55s, gave up, and had to abandon its
task. CLAUDE.md tells leads to prefer wait:false and to answer an ask with
bridge_send{turnId, content}; those two could not both be followed.

bridge_poll now returns a non-terminal ASKING phase carrying the question
and its turnId, and the ticket stays live so the worker's real reply still
lands on it. An unanswered ask returns the ticket to PENDING, because only
the question wait ended - the delegated turn continues. The 55s/115s ask
caps are unchanged: they exist because the worker's own MCP call would time
out, so widening them would only move the failure.

Two defects found reviewing the first revision, both from replacing
supplyAsync with a manually completed future:

- an exception inside the send left the future uncompleted, so the ticket
  stayed PENDING for the life of the daemon. Now caught and completed
  exceptionally.
- correlation was keyed by target, one entry per worker, registered before
  the session lock. With two tickets outstanding on one target the second
  overwrote the first, so a late reply could resolve the wrong ticket.
  Correlation is now per turn, the target entry exists only while that send
  owns the lock, and a reply with no live waiter still goes to the durable
  inbox as before.
2026-08-15 05:50:45 +02:00
Dai Ha 695da7418e Merge CB-573 (part 1): health classification model and the bridge_list capacity view
CI / contract (push) Successful in 45s
CI / build (push) Successful in 55s
Fleet capacity was invisible. A finished member held a terra slot until a
spawn was refused with 'at maxLoad: 2 live >= 2 cap', and nothing had told
the lead the slot was still held. bridge_list now reports, per configured
profile, maxLoad / live / free / reclaimable, and per member idleForSeconds
and reclaimable.

live comes from the same liveCountRef function placement consumes, so the
advertised free slots cannot drift from what bridge_spawn will accept. The
profile list is the union of configured and roster profiles: an empty
configured profile still appears with its full capacity, and a member whose
profile was removed from config stays visible rather than vanishing.

reclaimable is advisory. The bridge never spawns, stops or retasks a member
to improve utilisation: it has capacity facts but no work list, and choosing
work needs authority it does not have.

Also lands the pure health classifier, its precedence chain, the MUTE counter
and the pane budget. The classifier never reports IDLE while an accepted
delivery is open — IDLE is a claim that nothing is outstanding, and the
capacity view reads exactly that field.

Capacity dependencies are one required CapacitySource rather than defaulted
constructor arguments. A defaulted liveCount would report free slots that do
not exist, which is the dangerous direction; CapacitySource.none() omits the
block instead of inventing zeros.
2026-08-15 05:49:51 +02:00
Dai Ha abd26c796b CB-574: retain async task correlation
CI / build (pull_request) Failing after 1m14s
CI / contract (pull_request) Successful in 1m15s
2026-08-15 05:49:27 +02:00
Dai Ha 24559d81ac CB-573: require explicit capacity source
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Successful in 55s
2026-08-15 05:49:04 +02:00
Dai Ha 6c1c2c3994 CB-573: report empty configured profile capacity
CI / build (pull_request) Successful in 59s
CI / contract (pull_request) Successful in 1m0s
2026-08-15 05:45:39 +02:00
Dai Ha 01fab15713 CB-573: add fleet capacity view
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Successful in 1m25s
2026-08-15 05:43:46 +02:00
Dai Ha 0cd00e71c3 CB-574: surface async worker questions
CI / contract (pull_request) Successful in 42s
CI / build (pull_request) Successful in 56s
2026-08-15 05:43:43 +02:00
Dai Ha bf0ff2adbf CB-573: keep active delegations out of idle
CI / contract (pull_request) Successful in 45s
CI / build (pull_request) Successful in 1m27s
2026-08-15 05:40:37 +02:00
Dai Ha ed4bbc1c56 CB-573: add pure health classification model
CI / build (pull_request) Successful in 59s
CI / contract (pull_request) Successful in 1m15s
2026-08-15 05:37:30 +02:00
Dai Ha c456402cc5 Merge CB-572: reject a configured profile name as a send target
CI / build (push) Successful in 59s
CI / contract (push) Successful in 1m16s
A lead sent to sessionId "sol" — a profile name, not a terminal id. The
bridge accepted it, handed out a ticket, failed 60s later inside the
injector, and still reported the ticket as pending 20 minutes on. The
sender never learned anything and a whole delegation was lost.

bridge_send now rejects a target that exactly matches a configured
profile name, on both the blocking and the wait:false path, before any
ticket is issued. The error names the value and points at bridge_list.

The check is deliberately narrow. A target absent from the member roster
may still be a peer lead's terminal or a herdr-owned pane, so only a
value the bridge can prove is a profile is refused. profiles is a
required parameter on both send methods — the earlier revision kept
overloads that defaulted it to an empty set, which is the same silent
disable shape as CB-561.
2026-08-15 05:35:12 +02:00
Dai Ha 4f0bf667b1 CB-572: require profiles for send validation
CI / build (pull_request) Successful in 1m0s
CI / contract (pull_request) Successful in 1m2s
2026-08-15 05:33:15 +02:00
Dai Ha 61af9aa574 CB-572: reject profile names as send targets
CI / build (pull_request) Successful in 59s
CI / contract (pull_request) Successful in 1m20s
2026-08-15 05:29:37 +02:00
Dai Ha 3b59b34e76 CB-570 follow-up: rewrap an over-long javadoc line
CI / build (push) Successful in 58s
CI / contract (push) Successful in 1m17s
2026-08-15 05:12:58 +02:00
Dai Ha 65b38997f7 Merge CB-570: OpenCode receives the composed charter (PR #40)
One file, member-charter.md, not two. Two files would have made the U5
digest non-comparable between the Claude adapter and this one, which is
the whole point of the receipt; and nobody verified how OpenCode merges
multiple instruction files, so array order was an unverified dependency.

The OPENCODE_CONFIG condition widens to include a charter. It used to be
hasMcp() || hasCustomProvider(cfg), so a profile with a role charter but
no MCP and no custom provider would have got no config file and therefore
no charter — the feature silently doing nothing for that profile.

The file stays in the per-spawn temp dir, never the worktree: the
worktree is removed on release, the parity overlay already writes into
it, and CB-525's lesson was that config the bridge copied into a worktree
made a worker operate on the wrong tree. Being outside the repo is also
what stops it being committed, which a .gitignore line does not.
2026-08-15 05:12:26 +02:00
Dai Ha 7510f7649c Merge CB-569: Claude receives the composed charter (PR #39)
argvWithBridge used to gate the charter on cfg.hasMcp(), because the only
charter was the reply rule and telling a peer to call a tool it was not
given is a bug. A role charter is identity, not a tool instruction, so
the two gates are now separate: the MCP mount still depends on mcpUrl,
while --append-system-prompt depends only on the base having composed
something. A profile with a role charter and no MCP now gets its charter.

A null charter adds no flag at all. An empty --append-system-prompt is
not the same as no system prompt.
2026-08-15 05:11:39 +02:00
Dai Ha cea1183f75 CB-569: pass composed charter to Claude
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 50s
2026-08-15 05:09:30 +02:00
41 changed files with 2457 additions and 161 deletions
+10
View File
@@ -85,6 +85,16 @@ bind:
# backoffMs: 60000
# quietNudgeCap: 3
# Fleet health detection is dormant unless enabled. It reads one whole-fleet agent list per tick.
# It can run without a webhook; bridge_list then reports healthCoverage: detection-only.
# health:
# enabled: true
# intervalSeconds: 30 # minimum 15
# workingSuspectAfterSeconds: 600 # minimum 300
# paneProbeIntervalSeconds: 60 # minimum 60
# notifications:
# mode: disabled # disabled (default) or webhook
# herdr Unix socket. Omit to use the client default
# (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}).
herdrSocket: ~/.config/herdr/herdr.sock
@@ -23,6 +23,7 @@ import dev.ltms.bridged.mcp.BridgeMcp;
import dev.ltms.bridged.mcp.ConnectionIdentity;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.bridged.health.FleetHealthMonitor;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.mcp.LsofPeerPidLookup;
import dev.ltms.bridged.mcp.LsofProcessCwdLookup;
@@ -275,9 +276,9 @@ public final class Bridged {
}
@Override
public void onDelivered(String target) {
completion.onDelivered(target);
sessions.onDelivered(target);
public void onDelivered(String target, dev.ltms.bridged.msg.TurnToken token) {
completion.onDelivered(target, token);
sessions.onDelivered(target, token);
}
@Override
@@ -285,6 +286,12 @@ public final class Bridged {
completion.onTurnFailed(target);
sessions.onTurnFailed(target);
}
@Override
public void onTurnFailed(String target, String reason) {
completion.onTurnFailed(target, reason);
sessions.onTurnFailed(target);
}
};
Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads),
presence::forget);
@@ -348,6 +355,26 @@ public final class Bridged {
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
pushLoop, metrics);
// Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller.
final FleetHealthMonitor healthMonitor;
var healthScheduler = Executors.newSingleThreadScheduledExecutor(r ->
Thread.ofVirtual().name("bridge-health-").unstarted(r));
if (cfg.health() != null && cfg.health().isEnabled()) {
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
System::nanoTime, cfg.health().intervalOrDefault());
String coverage = FleetHealthMonitor.coverage(true,
cfg.health().notifications() != null && cfg.health().notifications().configured());
if ("detection-only".equals(coverage)) {
log.warn("fleet health: {} (no notification sink configured)", coverage);
} else {
log.info("fleet health: {}", coverage);
}
healthMonitor.start();
} else {
healthMonitor = null;
healthScheduler.shutdownNow();
}
// CB-520: the reply inbox only consumes for agents this gateway owns. own on acquire,
// release on teardown. Do this before CB-516 so the inbox is owned before any reply can land.
sessions.onAcquire(replyInbox::own);
@@ -383,7 +410,16 @@ public final class Bridged {
}
BridgeMcp mcp = new BridgeMcp(messages, workers, sessions, identity, presence,
primaryRegistry, callers, metrics);
primaryRegistry, callers, metrics, new BridgeMcp.CapacitySource(profile -> liveCountRef.get().apply(profile),
profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.maxLoad();
}, () -> config.get().profiles().keySet(), System::nanoTime),
new BridgeMcp.HealthCoverageSource(() -> {
var health = config.get().health();
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
health != null && health.notifications() != null && health.notifications().configured());
}));
// CB-559: opt-in config reload. With no `configReload:` block nothing is constructed, so an
// upgraded daemon behaves exactly as before — the file is read once at boot and never again.
@@ -404,6 +440,7 @@ public final class Bridged {
messages.close();
pushLoop.close();
if (heartbeat != null) heartbeat.close(); // CB-551: stop the idle-lead heartbeat scheduler
if (healthMonitor != null) healthMonitor.stop();
if (configWatcher != null) configWatcher.stop(); // CB-559: stop polling the config file
mcp.close();
if (reaper != null) reaper.stop();
@@ -73,6 +73,7 @@ public record BridgedConfig(
Primary primary,
Fleet fleet,
LeadHeartbeat leadHeartbeat,
Health health,
String placement,
Auth auth,
ConfigReload configReload) {
@@ -83,7 +84,16 @@ public record BridgedConfig(
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, String placement, Auth auth) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, placement, auth, null);
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, null);
}
/** Back-compat form before the optional {@code health:} block was added. */
public BridgedConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, String placement, Auth auth, ConfigReload configReload) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, null, placement, auth, configReload);
}
/**
@@ -377,6 +387,17 @@ public record BridgedConfig(
boolean clearAfterTurn) {
}
/** Optional fleet detection. A missing block stays dormant. */
@JsonIgnoreProperties(ignoreUnknown = true)
public record Health(Boolean enabled, Integer intervalSeconds, Integer workingSuspectAfterSeconds,
Integer paneProbeIntervalSeconds, Notifications notifications) {
public boolean isEnabled() { return Boolean.TRUE.equals(enabled); }
public int intervalOrDefault() { return Math.max(15, intervalSeconds == null ? 30 : intervalSeconds); }
public record Notifications(String mode) {
public boolean configured() { return "webhook".equalsIgnoreCase(mode); }
}
}
/**
* External AMQP broker for durable, cross-restart reply delivery (CB-307 Stage 2). Its mere
* presence swaps the in-memory {@code ReplyInbox} for the AMQP-backed adapter; absent, bridged
@@ -812,7 +833,7 @@ public record BridgedConfig(
private static final Set<String> KNOWN_TOP_LEVEL_KEYS = Set.of(
"bind", "herdrSocket", "profiles", "guard", "worktreeRoot",
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "placement", "auth", "configReload");
"leadHeartbeat", "health", "placement", "auth", "configReload");
/** Load and validate config from {@code path}. */
public static BridgedConfig load(Path path) {
@@ -1082,7 +1103,7 @@ public record BridgedConfig(
// defaults the fields of a block that IS present. Defaulting it here would start watching
// the file for every config that never asked to be watched.
return new BridgedConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
broker, primary, f, leadHeartbeat, placementOrDefault, a, configReload);
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload);
}
/**
@@ -0,0 +1,40 @@
package dev.ltms.bridged.health;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.session.MemberSession;
/**
* Pure classifier. Collection and repair are deliberately outside this package.
* {@link HealthState#ERROR_ON_SCREEN} is not decided yet because it needs a bounded pane detection
* read and an adapter-specific fatal signature; status facts alone must not guess it.
*/
public final class FleetHealth {
private FleetHealth() { }
public static HealthDecision decide(HealthSnapshot s, HealthPrior prior, long nowNanos) {
if (s.controlLinkDown()) return result(HealthState.CONTROL_LINK_DOWN, false);
if (s.targetNotFound()) return result(HealthState.GONE, false);
if (s.sessionState() == MemberSession.State.SPAWNING && !s.present() && s.readinessGraceElapsed()) {
return result(HealthState.NEVER_READY, false);
}
if (s.orphanedDelegation()) return result(HealthState.DELEGATION_ORPHANED, false);
boolean disagreement = s.sessionState() == MemberSession.State.BUSY && s.acceptedDelivery()
&& (s.liveStatus() == AgentStatus.IDLE || s.liveStatus() == AgentStatus.DONE);
if (disagreement && prior.busyButDone()) return result(HealthState.TURN_BOUNDARY_LOST, true);
if (s.stalled()) return result(HealthState.STALL_SUSPECTED, disagreement);
if (s.replyStranded()) return result(HealthState.REPLY_STRANDED, disagreement);
if (s.queuedDelivery() || s.inboxMessage()) return result(HealthState.WORK_PENDING, disagreement);
if (s.sessionState() == MemberSession.State.SPAWNING) return result(HealthState.STARTING, disagreement);
if (s.acceptedDelivery() && s.liveStatus() == AgentStatus.BLOCKED) {
return result(HealthState.BLOCKED_AMBIGUOUS, disagreement);
}
// An accepted delivery remains bridge work even when herdr is late, unknown, or has already
// reported DONE once. It cannot be IDLE until the delegation has resolved.
if (s.acceptedDelivery()) return result(HealthState.WORKING, disagreement);
return result(HealthState.IDLE, disagreement);
}
private static HealthDecision result(HealthState state, boolean disagreement) {
return new HealthDecision(state, new HealthPrior(disagreement));
}
}
@@ -0,0 +1,107 @@
package dev.ltms.bridged.health;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.session.MemberSession;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.function.LongSupplier;
import java.util.function.Supplier;
/** Slow whole-fleet evidence collection. It is deliberately separate from the delivery poller. */
public final class FleetHealthMonitor {
private static final Logger log = LoggerFactory.getLogger(FleetHealthMonitor.class);
private final AgentControl agents;
private final Supplier<List<MemberSession>> roster;
private final MessageService messages;
private final ScheduledExecutorService scheduler;
private final LongSupplier clock;
private final long intervalSeconds;
private final Map<String, HealthPrior> priors = new HashMap<>();
private final Map<String, HealthState> states = new HashMap<>();
// These facts need the evidence publishers introduced by later M4 units. They are not negatives.
private static final boolean NOT_YET_OBSERVED = false;
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds) {
this.agents = agents;
this.roster = roster;
this.messages = messages;
this.scheduler = scheduler;
this.clock = clock;
this.intervalSeconds = intervalSeconds;
}
/** Pure per-member decision seam. */
static HealthDecision decide(HealthSnapshot snapshot, HealthPrior prior, long nowNanos) {
return FleetHealth.decide(snapshot, prior, nowNanos);
}
public void start() { scheduler.schedule(this::tick, intervalSeconds, TimeUnit.SECONDS); }
public void stop() { scheduler.shutdownNow(); }
// Package-private so tests can run one tick without waiting.
void tick() {
try {
List<Agent> agentsNow = agents.list(); // Exactly one list call for this complete observation.
List<MemberSession> rosterNow = roster.get(); // One in-memory roster snapshot for this tick.
Map<String, Agent> live = new HashMap<>();
for (Agent agent : agentsNow) live.put(agent.terminalId(), agent);
HashSet<String> current = new HashSet<>();
for (MemberSession session : rosterNow) {
current.add(session.terminalId());
Agent agent = live.get(session.terminalId());
AgentStatus status = agent == null ? AgentStatus.UNKNOWN : agent.status();
boolean accepted = messages.hasAcceptedDelivery(session.terminalId());
HealthSnapshot snapshot = new HealthSnapshot(session.state(), status, accepted, NOT_YET_OBSERVED,
messages.hasInboxMessage(session.terminalId()), agent != null, NOT_YET_OBSERVED,
NOT_YET_OBSERVED, NOT_YET_OBSERVED, NOT_YET_OBSERVED, NOT_YET_OBSERVED, NOT_YET_OBSERVED);
HealthDecision decision = decide(snapshot, priors.getOrDefault(session.terminalId(), HealthPrior.NONE),
clock.getAsLong());
priors.put(session.terminalId(), decision.prior());
reportTransition(session.terminalId(), decision.state());
}
priors.keySet().retainAll(current);
states.keySet().retainAll(current);
} catch (Throwable error) {
// A list failure is health evidence, and must never kill the monitor's only scheduler task.
log.warn("fleet health collection failed; will retry next tick", error);
} finally {
if (!scheduler.isShutdown()) {
scheduler.schedule(this::tick, intervalSeconds, TimeUnit.SECONDS);
}
}
}
void reportTransition(String target, HealthState next) {
HealthState previous = states.put(target, next);
if (previous == next) return;
if (fault(next)) {
log.warn("fleet health member={} state={} previous={}", target, next, previous);
} else if (previous != null && fault(previous)) {
log.info("fleet health member={} recovered state={} previous={}", target, next, previous);
}
}
private static boolean fault(HealthState state) {
return switch (state) {
case NEVER_READY, GONE, TURN_BOUNDARY_LOST, ERROR_ON_SCREEN, STALL_SUSPECTED,
MUTE, REPLY_STRANDED, DELEGATION_ORPHANED, CONTROL_LINK_DOWN -> true;
default -> false;
};
}
public static String coverage(boolean enabled, boolean notificationConfigured) {
return !enabled ? "off" : notificationConfigured ? "full" : "detection-only";
}
}
@@ -0,0 +1,4 @@
package dev.ltms.bridged.health;
/** Classification plus the private fact that the next pure decision needs. */
public record HealthDecision(HealthState state, HealthPrior prior) { }
@@ -0,0 +1,6 @@
package dev.ltms.bridged.health;
/** Private cross-tick observation. It is deliberately not a reported health value. */
public record HealthPrior(boolean busyButDone) {
public static final HealthPrior NONE = new HealthPrior(false);
}
@@ -0,0 +1,11 @@
package dev.ltms.bridged.health;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.session.MemberSession;
/** Read-only facts from one fleet collection tick. */
public record HealthSnapshot(MemberSession.State sessionState, AgentStatus liveStatus,
boolean acceptedDelivery, boolean queuedDelivery, boolean inboxMessage,
boolean present, boolean targetNotFound, boolean controlLinkDown,
boolean readinessGraceElapsed, boolean orphanedDelegation,
boolean replyStranded, boolean stalled) { }
@@ -0,0 +1,8 @@
package dev.ltms.bridged.health;
/** Health classifications reported for a member. */
public enum HealthState {
STARTING, IDLE, WORKING, WORK_PENDING, BLOCKED_AMBIGUOUS,
NEVER_READY, GONE, TURN_BOUNDARY_LOST, ERROR_ON_SCREEN, STALL_SUSPECTED,
MUTE, REPLY_STRANDED, DELEGATION_ORPHANED, CONTROL_LINK_DOWN
}
@@ -0,0 +1,25 @@
package dev.ltms.bridged.health;
import dev.ltms.bridged.msg.Rendezvous;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* Counts turns that ended via the completion fallback instead of {@code bridge_reply}.
* MUTE is an observation by target and profile, not a classifier state and never suppresses faults.
*/
public final class MuteCounter {
private final Map<String, Integer> byTarget = new ConcurrentHashMap<>();
private final Map<String, Integer> byProfile = new ConcurrentHashMap<>();
/** Record only fallback completion; a structured reply does not make a member mute. */
public void observe(String target, String profile, Rendezvous.Kind kind) {
if (kind != Rendezvous.Kind.COMPLETION) return;
byTarget.merge(target, 1, Integer::sum);
byProfile.merge(profile, 1, Integer::sum);
}
public int forTarget(String target) { return byTarget.getOrDefault(target, 0); }
public int forProfile(String profile) { return byProfile.getOrDefault(profile, 0); }
}
@@ -0,0 +1,26 @@
package dev.ltms.bridged.health;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
/** Fixed pane-probe limits. Pane content is never retained here. */
public final class PaneBudget {
public static final long COOLDOWN_NANOS = 60_000_000_000L;
public static final int MAX_PER_TICK = 2;
private final Map<String, Long> lastProbe = new HashMap<>();
private int cursor;
public List<String> choose(List<String> candidates, long nowNanos, long configuredCooldownNanos) {
long cooldown = Math.max(COOLDOWN_NANOS, configuredCooldownNanos);
List<String> out = new ArrayList<>();
for (int n = 0; n < candidates.size() && out.size() < MAX_PER_TICK; n++) {
String target = candidates.get((cursor + n) % candidates.size());
Long last = lastProbe.get(target);
if (last == null || nowNanos - last >= cooldown) { out.add(target); lastProbe.put(target, nowNanos); }
}
if (!candidates.isEmpty()) cursor = (cursor + 1) % candidates.size();
return List.copyOf(out);
}
}
@@ -2,6 +2,7 @@ package dev.ltms.bridged.inject;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.msg.TurnToken;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -83,17 +84,17 @@ public final class CompletionResolver implements TurnListener {
}
@Override
public void onDelivered(String target) {
public void onDelivered(String target, TurnToken token) {
// Capture the exact waiter this turn belongs to (CB-116) and snapshot the pane's pre-turn
// content — what it shows *before* the just-delivered turn produces output — as the staleness
// reference (CB-115). Done synchronously (like the delivering send itself) so both are in
// place before this turn's completion can fire.
captureBaseline(target);
captureBaseline(target, token);
}
/** Capture the in-flight turn: its waiter and pre-turn baseline (the testable core of {@link #onDelivered}). */
void captureBaseline(String target) {
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
void captureBaseline(String target, TurnToken token) {
CompletableFuture<Rendezvous.Resolution> waiter = token.waiter();
if (waiter == null) {
inFlight.remove(target); // no send is waiting on this delivery — nothing to resolve later
return;
@@ -137,7 +138,13 @@ public final class CompletionResolver implements TurnListener {
@Override
public void onTurnFailed(String target) {
InFlight turn = inFlight.get(target);
Thread.ofVirtual().name("turn-failed-" + target).start(() -> fail(target, turn));
Thread.ofVirtual().name("turn-failed-" + target).start(() -> fail(target, turn, null));
}
@Override
public void onTurnFailed(String target, String reason) {
InFlight turn = inFlight.get(target);
Thread.ofVirtual().name("turn-failed-" + target).start(() -> fail(target, turn, reason));
}
/** Synchronous resolve (the unit-testable core of {@link #onTurnComplete}). */
@@ -192,6 +199,11 @@ public final class CompletionResolver implements TurnListener {
/** Synchronous fail (the unit-testable core of {@link #onTurnFailed}). */
void fail(String target, InFlight turn) {
fail(target, turn, null);
}
/** Synchronous fail with an optional reason supplied by a dropped worker queue. */
void fail(String target, InFlight turn, String explicitReason) {
// A never-delivered readiness failure has no in-flight record but still has a blocked send;
// fall back to the currently-registered waiter (unambiguous — that send never completed, so
// no next turn exists to confuse it with).
@@ -201,16 +213,18 @@ public final class CompletionResolver implements TurnListener {
inFlight.remove(target, turn); // nobody blocked on this worker — nothing to fail
return;
}
String reason;
try {
reason = clip(agents.read(target, SCRAPE_SOURCE));
} catch (RuntimeException e) {
reason = "";
}
if (reason.isBlank()) {
// No screen to scrape — either the worker is stuck (CB-109) or gone (CB-110).
reason = "worker did not reply; its turn ended in an unrecoverable state "
+ "(worker unreachable or stuck)";
String reason = explicitReason;
if (reason == null || reason.isBlank()) {
try {
reason = clip(agents.read(target, SCRAPE_SOURCE));
} catch (RuntimeException e) {
reason = "";
}
if (reason.isBlank()) {
// No screen to scrape — either the worker is stuck (CB-109) or gone (CB-110).
reason = "worker did not reply; its turn ended in an unrecoverable state "
+ "(worker unreachable or stuck)";
}
}
if (rendezvous.resolveFailure(waiter, reason)) {
inFlight.remove(target, turn);
@@ -2,6 +2,7 @@ package dev.ltms.bridged.inject;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.msg.TurnToken;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -129,7 +130,7 @@ public final class Injector {
}
/** A pending message and the future that completes when it has been delivered. */
private record Pending(String text, CompletableFuture<Void> delivered) {
private record Pending(String text, TurnToken token, CompletableFuture<Void> delivered) {
}
/** Per-worker delivery state, guarded by its own monitor (single writer per worker). */
@@ -159,9 +160,9 @@ public final class Injector {
* <p>Uses an atomic map update so a concurrent {@link #drop} cannot slip between "find the
* target" and "queue the message" and orphan it in a target it just removed.
*/
public CompletableFuture<Void> enqueue(String target, String text) {
public CompletableFuture<Void> enqueue(String target, String text, TurnToken token) {
CompletableFuture<Void> delivered = new CompletableFuture<>();
Pending p = new Pending(text, delivered);
Pending p = new Pending(text, token, delivered);
targets.compute(target, (_, existing) -> {
Target t = (existing != null) ? existing : new Target();
t.add(p); // synchronized on the Target monitor — atomic with a concurrent drop
@@ -356,7 +357,7 @@ public final class Injector {
} else {
// Baseline the pane's pre-turn content so a misattributed completion (no new output)
// can't resolve this send with the previous turn's stale answer (CB-115).
turnListener.onDelivered(target);
turnListener.onDelivered(target, sent.token());
sent.delivered().complete(null);
}
}
@@ -410,8 +411,8 @@ public final class Injector {
for (Pending p : pending) {
p.delivered().completeExceptionally(cause);
}
if (hadDeliveredTurn) {
turnListener.onTurnFailed(target);
}
// A queued send has no in-flight record, while a delivered turn does. CompletionResolver
// handles both forms and resolves its waiter at most once.
turnListener.onTurnFailed(target, cause.getMessage());
}
}
@@ -1,5 +1,7 @@
package dev.ltms.bridged.inject;
import dev.ltms.bridged.msg.TurnToken;
/**
* Notified when a worker's delegated turn is observed to complete — a confirmed
* {@code WORKING → IDLE} transition after a delivery. This is the CB-106 completion signal the
@@ -41,6 +43,14 @@ public interface TurnListener {
default void onTurnFailed(String target) {
}
/**
* As {@link #onTurnFailed(String)}, carrying the reason a worker became unreachable. The default
* keeps existing listeners working while allowing the completion resolver to report a useful cause.
*/
default void onTurnFailed(String target, String reason) {
onTurnFailed(target);
}
/**
* A message was just delivered into {@code target}'s pane (CB-115). Fired so the completion
* resolver can snapshot the pane's pre-turn content: a later {@link #onTurnComplete} whose
@@ -49,7 +59,7 @@ public interface TurnListener {
* resolve the send with the previous turn's stale answer. A default no-op keeps the interface
* functional for callers that don't scrape.
*/
default void onDelivered(String target) {
default void onDelivered(String target, TurnToken token) {
}
/** No-op default for callers that only need delivery, not completion signalling. */
@@ -0,0 +1,30 @@
package dev.ltms.bridged.logging;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.turbo.TurboFilter;
import ch.qos.logback.core.spi.FilterReply;
import io.modelcontextprotocol.spec.McpSchema;
import org.slf4j.Marker;
/** Suppresses only the SDK warning for the normal MCP cancellation notification. */
public final class McpCancelledNotificationFilter extends TurboFilter {
static final String LOGGER = "io.modelcontextprotocol.spec.McpStreamableServerSession";
static final String UNHANDLED_NOTIFICATION = "No handler registered for notification method: {}";
@Override
public FilterReply decide(Marker marker, Logger logger, Level level, String format, Object[] params,
Throwable throwable) {
if (level == Level.WARN
&& LOGGER.equals(logger.getName())
&& UNHANDLED_NOTIFICATION.equals(format)
&& params != null
&& params.length == 1
&& params[0] instanceof McpSchema.JSONRPCNotification notification
&& "notifications/cancelled".equals(notification.method())) {
return FilterReply.DENY;
}
return FilterReply.NEUTRAL;
}
}
@@ -33,7 +33,10 @@ import jakarta.servlet.http.HttpServlet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.function.Function;
import java.util.function.LongSupplier;
import java.util.function.Supplier;
import java.util.stream.Collectors;
/**
@@ -73,17 +76,20 @@ public final class BridgeMcp {
private final McpSyncServer server;
private final CallerResolver authz; // CB-501: null → authorization not enforced (legacy)
private final Metrics metrics; // CB-502: null → auth failures not counted
private final CapacitySource capacity;
private final HealthCoverageSource healthCoverage;
/**
* Legacy constructor — no authorization. Retained so existing tests exercise tool behaviour
* without an auth fixture.
*/
public BridgeMcp(MessageService messages, PeerLauncher workers,
SessionManager sessions, ConnectionIdentity identity, MemberPresence presence,
PrimaryRegistry primaryRegistry) {
this(messages, workers, sessions, identity, presence, primaryRegistry, null, null);
/** Capacity facts used by {@code bridge_list}; production must supply the placement live count. */
public record CapacitySource(Function<String, Integer> liveCount, Function<String, Integer> maxLoad,
Supplier<Set<String>> configuredProfiles, LongSupplier clock) {
/** Inert test-only source. It omits capacity rather than inventing zero live counts. */
public static CapacitySource none() { return new CapacitySource(_ -> 0, _ -> null, Set::of, System::nanoTime); }
boolean available() { return !configuredProfiles.get().isEmpty(); }
}
/** Coverage is supplied by the health wiring, not inferred from a missing dependency. */
public record HealthCoverageSource(Supplier<String> value) { }
/**
* @param callers resolves each call's {@link Principal}; {@code null} disables authorization.
* This surface needs its own enforcement: {@code /mcp} is a raw servlet on
@@ -91,9 +97,11 @@ public final class BridgeMcp {
* filter, so the REST guard does not cover it.
* @param metrics registry for auth-failure counting; may be {@code null}
*/
public BridgeMcp(MessageService messages, PeerLauncher workers,
SessionManager sessions, ConnectionIdentity identity, MemberPresence presence,
PrimaryRegistry primaryRegistry, CallerResolver callers, Metrics metrics) {
public BridgeMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry,
CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage) {
this.capacity = capacity;
this.healthCoverage = healthCoverage;
McpJsonMapper json = new JacksonMcpJsonMapperSupplier().get();
this.transport = HttpServletStreamableServerTransportProvider.builder()
.jsonMapper(json)
@@ -151,8 +159,8 @@ public final class BridgeMcp {
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, caller);
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
return Boolean.FALSE.equals(a.get("wait"))
? sendAsync(messages, target, content, onAccepted)
: send(messages, target, content, timeoutMs(a), onAccepted);
? sendAsync(messages, target, content, onAccepted, workers.profiles())
: send(messages, target, content, timeoutMs(a), onAccepted, workers.profiles());
})
// bridge_reply's identity is the CONNECTION, never an argument — so the authz check
// is "is this caller a worker at all", and it can only ever reply as itself.
@@ -207,7 +215,7 @@ public final class BridgeMcp {
.toolCall(listTool(), (exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return listFleet(workers, sessions,
return listFleet(workers, sessions, messages, capacity, healthCoverage,
callers == null ? Map.of() : callers.leads(),
callerTerminal(exchange));
})
@@ -379,21 +387,22 @@ public final class BridgeMcp {
// --- tool logic (thin adapters over the services; unit-testable) ---------------------------
/** {@code bridge_send}: delegate {@code content} to a worker session and block for its reply. */
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content, Long timeoutMs) {
return send(messages, sessionId, content, timeoutMs, null);
}
/**
* As {@link #send(MessageService, String, String, Long)}, wiring an accepted-delivery hook
* {@code bridge_send}: delegate {@code content} to a worker session and block for its reply.
* The configured profiles are required so a profile name can never bypass target validation.
*
* (CB-548): {@code onAccepted} records delegator ownership the instant the send is accepted, so
* a BUSY interloper never claims a turn it did not win. {@code null} disables recording.
*/
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content,
Long timeoutMs, Runnable onAccepted) {
Long timeoutMs, Runnable onAccepted, Set<String> profiles) {
if (isBlank(sessionId) || isBlank(content)) {
return error("sessionId and content are required");
}
McpSchema.CallToolResult targetError = profileTargetError(sessionId, profiles);
if (targetError != null) {
return targetError;
}
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
try {
return formatReply(messages.send(sessionId, content, timeout, onAccepted), timeout);
@@ -463,24 +472,33 @@ public final class BridgeMcp {
/**
* {@code bridge_send} with {@code wait:false}: delegate {@code content} and return a ticket
* immediately (fire-and-poll), so a long task isn't cut off by the caller's MCP call timeout.
*/
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content) {
return sendAsync(messages, sessionId, content, null);
}
/**
* As {@link #sendAsync(MessageService, String, String)}, wiring the accepted-delivery hook
* The configured profiles are required so a profile name can never bypass target validation.
*
* This wires the accepted-delivery hook
* (CB-548) so an async flooding send records delegator ownership exactly once it is accepted.
*/
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
Runnable onAccepted) {
Runnable onAccepted, Set<String> profiles) {
if (isBlank(sessionId) || isBlank(content)) {
return error("sessionId and content are required");
}
McpSchema.CallToolResult targetError = profileTargetError(sessionId, profiles);
if (targetError != null) {
return targetError;
}
String ticket = messages.sendAsync(sessionId, content, onAccepted);
return text("accepted — task delegated. Poll bridge_poll with ticket=" + ticket);
}
/** A configured profile is never a send target; other unknown values may be herdr-owned panes. */
private static McpSchema.CallToolResult profileTargetError(String sessionId, Set<String> profiles) {
if (profiles.contains(sessionId)) {
return error("unknown send target \"" + sessionId + "\": it is a configured profile name, not a "
+ "session id. Call bridge_list to find a member or lead sessionId.");
}
return null;
}
/** {@code bridge_poll}: check an async delegation by ticket, or drain a worker's inbox by target. */
static McpSchema.CallToolResult poll(MessageService messages, String ticket, String target) {
if (!isBlank(target)) {
@@ -502,6 +520,9 @@ public final class BridgeMcp {
? "[done — worker finished without a structured bridge_reply; transcript tail follows]\n" + v.reply()
: v.reply());
case PENDING -> text("[pending — " + v.detail() + "]");
case ASKING -> text("[question — worker is waiting for your answer]\n" + v.reply()
+ "\n\nAnswer it by calling bridge_send again with turnId=\"" + v.turnId()
+ "\" and content set to your answer; the worker resumes the same turn.");
case FAILED -> text("[failed — " + v.detail() + "]");
};
}
@@ -708,6 +729,12 @@ public final class BridgeMcp {
*/
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions,
Map<String, String> leads, String selfTerm) {
return listFleet(workers, sessions, null, CapacitySource.none(), new HealthCoverageSource(() -> "off"), leads, selfTerm);
}
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
Map<String, String> leads, String selfTerm) {
try {
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
@@ -717,15 +744,57 @@ public final class BridgeMcp {
.sorted(Map.Entry.comparingByValue())
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm))
.toList();
List<Map<String, Object>> out = sessions.roster().stream()
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
List<MemberSession> roster = sessions.roster();
List<Map<String, Object>> out = roster.stream()
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
.toList();
return text(json(Map.of("leads", leadRows, "members", out)));
Set<String> profiles = new java.util.TreeSet<>(capacity.configuredProfiles().get());
roster.stream().map(MemberSession::profile).forEach(profiles::add);
Map<String, Object> result = new LinkedHashMap<>();
result.put("leads", leadRows); result.put("members", out);
result.put("healthCoverage", healthCoverage.value().get());
if (capacity.available()) result.put("capacity", profiles.stream()
.map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages,
capacity.clock().getAsLong())).toList());
return text(json(result));
} catch (HerdrException e) {
return error("herdr error listing the fleet: " + e.getMessage());
}
}
/**
* Capacity is advisory only. {@code reclaimable} says there is no bridge work, not that bridged
* may stop the member: the bridge has capacity facts but no work list, and choosing work needs
* authority it does not have. {@code idleForSeconds} is derived from monotonic nanoTime and has
* no wall-clock meaning across a daemon restart.
*/
private static Map<String, Object> memberCapacityView(MemberSession session, Agent live,
MessageService messages, long nowNanos) {
Map<String, Object> row = SessionManager.rosterView(session, live);
boolean open = messages != null && messages.hasAcceptedDelivery(session.terminalId());
boolean inbox = messages != null && messages.hasInboxMessage(session.terminalId());
boolean reclaimable = (session.state() == MemberSession.State.READY || session.state() == MemberSession.State.DONE)
&& !open && !inbox;
row.put("reclaimable", reclaimable);
row.put("idleForSeconds", reclaimable ? Math.max(0, (nowNanos - session.lastActivityAtNanos()) / 1_000_000_000L) : null);
return row;
}
private static Map<String, Object> capacityView(String profile, Function<String, Integer> liveCount,
Function<String, Integer> maxLoad, List<MemberSession> roster,
MessageService messages, long nowNanos) {
Integer cap = maxLoad.apply(profile);
int live = liveCount.apply(profile);
int reclaimable = (int) roster.stream().filter(s -> profile.equals(s.profile()))
.filter(s -> (s.state() == MemberSession.State.READY || s.state() == MemberSession.State.DONE))
.filter(s -> messages == null || (!messages.hasAcceptedDelivery(s.terminalId()) && !messages.hasInboxMessage(s.terminalId())))
.count();
Map<String, Object> row = new LinkedHashMap<>();
row.put("profile", profile); row.put("maxLoad", cap); row.put("live", live);
row.put("free", cap == null ? null : Math.max(0, cap - live)); row.put("reclaimable", reclaimable);
return row;
}
/**
* One lead's row: its address, its name, and whether it can be reached right now.
*
@@ -181,8 +181,8 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
// id via -r and passes no --session-id (the two conflict). Both are injected before the
// model flag so --model keeps outranking the operator's own argv.
// mutableArgv: argvWithBridge may hand back the profile's own (immutable) List.of when it
// has no MCP — session flags must be added into a list we own.
List<String> argv = mutableArgv(argvWithBridge(cfg));
// has neither MCP nor a charter — session flags must be added into a list we own.
List<String> argv = mutableArgv(argvWithBridge(cfg, spec.charter()));
String agentSessionId = applySessionIdentity(argv, spec.sessionName(), spec.resumeSessionId());
return new Launch(workerEnv, argvWithModel(argv, cfg), agentSessionId);
}
@@ -217,22 +217,26 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
}
/**
* The launch argv, plus — when {@code worker.mcpUrl} is set — inline {@code --mcp-config} for
* the bridge server and {@code --append-system-prompt} for the {@link #REPLY_CHARTER}. Neither
* touches the profile's config; both are pure command-line flags. This inline-flag mount is
* Claude Code specific — other adapters mount MCP and instructions their own way.
* The launch argv, plus an inline {@code --mcp-config} when {@code worker.mcpUrl} is set and
* {@code --append-system-prompt} when the base composed a charter. Neither touches the profile's
* config; both are pure command-line flags. This inline-flag mount is Claude Code specific —
* other adapters mount MCP and instructions their own way.
*/
private List<String> argvWithBridge(BridgedConfig.Profile cfg) {
if (!cfg.hasMcp()) {
private List<String> argvWithBridge(BridgedConfig.Profile cfg, String charter) {
if (!cfg.hasMcp() && charter == null) {
return cfg.argv();
}
String mcpJson = "{\"mcpServers\":{\"bridge\":{\"type\":\"http\",\"url\":\""
+ cfg.mcpUrl() + "\"}}}";
List<String> argv = mutableArgv(cfg.argv());
argv.add("--mcp-config");
argv.add(mcpJson);
argv.add("--append-system-prompt");
argv.add(REPLY_CHARTER);
if (cfg.hasMcp()) {
String mcpJson = "{\"mcpServers\":{\"bridge\":{\"type\":\"http\",\"url\":\""
+ cfg.mcpUrl() + "\"}}}";
argv.add("--mcp-config");
argv.add(mcpJson);
}
if (charter != null) {
argv.add("--append-system-prompt");
argv.add(charter);
}
return argv;
}
@@ -163,10 +163,10 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* {@inheritDoc}
*
* <p>Builds the opencode launch: no {@code ANTHROPIC_*} and no guard (opencode reads its own
* provider credentials); when the profile mounts the bridge MCP or has a member charter, generate an ephemeral
* {@code opencode.json} (remote MCP server + member-charter instructions) and point the worker at
* it via {@code OPENCODE_CONFIG}; carry the parity-neutral git-forge grant; and select the model
* with {@code -m}.
* provider credentials); when the profile mounts the bridge MCP or has a member charter,
* generate an ephemeral {@code opencode.json} (remote MCP server + member-charter instructions)
* and point the worker at it via {@code OPENCODE_CONFIG}; carry the parity-neutral git-forge
* grant; and select the model with {@code -m}.
*/
@Override
protected Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec) {
@@ -127,6 +127,8 @@ public final class MessageService {
public enum Phase {
/** Delegated and in flight — queued for the worker or being worked. */
PENDING,
/** The worker is paused in {@code bridge_ask}; {@link TaskView#reply} and {@link TaskView#turnId} identify it. */
ASKING,
/** The worker's turn finished; {@link TaskView#reply} holds the answer. */
DONE,
/** The delegation could not complete (timed out, worker gone, or busy). */
@@ -136,16 +138,28 @@ public final class MessageService {
/**
* A poll snapshot of an async delegation.
*
* @param reply the answer when {@link #phase} is {@link Phase#DONE}, else {@code null}
* @param reply the answer when {@link #phase} is {@link Phase#DONE}, or the question when
* {@link #phase} is {@link Phase#ASKING}; otherwise {@code null}
* @param replySource {@code "reply"} (structured {@code bridge_reply}) or {@code "transcript"}
* (completion scrape) when {@link Phase#DONE}, else {@code null}
* @param detail a human note (live worker status while pending, or the failure reason)
* @param detail a human note (live worker status while pending, ask state, or failure reason)
* @param turnId correlation id for an {@link Phase#ASKING} ticket, else {@code null}
*/
public record TaskView(String ticket, Phase phase, String reply, String replySource, String detail) {
public record TaskView(String ticket, Phase phase, String reply, String replySource, String detail,
String turnId) {
}
/** An in-flight or finished async delegation, keyed by its ticket. */
private record Task(String target, CompletableFuture<Reply> future, long createdNanos) {
private static final class Task {
private final String target;
private final CompletableFuture<Reply> future = new CompletableFuture<>();
private final long createdNanos = System.nanoTime();
private volatile Reply question;
private volatile String turnId;
private Task(String target) {
this.target = target;
}
}
private final AgentControl agents;
@@ -156,6 +170,11 @@ public final class MessageService {
private final Metrics metrics; // CB-502: nullable — no registry in unit tests
private final ConcurrentHashMap<String, ReentrantLock> sessionLocks = new ConcurrentHashMap<>();
private final ConcurrentHashMap<String, Task> tasks = new ConcurrentHashMap<>();
/** Async task that owns each exact forward rendezvous waiter. */
private final ConcurrentHashMap<CompletableFuture<Rendezvous.Resolution>, Task> asyncTasksByWaiter =
new ConcurrentHashMap<>();
/** Async tickets paused on a specific {@code bridge_ask} turn. */
private final ConcurrentHashMap<String, Task> asyncTasksByTurn = new ConcurrentHashMap<>();
private final AtomicLong ticketSeq = new AtomicLong();
private final ExecutorService asyncExecutor = Executors.newThreadPerTaskExecutor(
Thread.ofVirtual().name("bridge-async-", 0).factory());
@@ -202,6 +221,16 @@ public final class MessageService {
return agents.status(target);
}
/** Read-only delegation fact for fleet views. */
public boolean hasAcceptedDelivery(String target) {
return rendezvous.isWaiting(target);
}
/** Read-only inbox fact for fleet views. */
public boolean hasInboxMessage(String target) {
return !inbox.peek(target).isEmpty();
}
/**
* Route a worker's explicit {@code bridge_reply}: resolve an open send, or queue it in the
* inbox if no send is currently open. Unlike the bare {@link Rendezvous#resolve}, a no-waiter
@@ -272,14 +301,18 @@ public final class MessageService {
*/
public boolean abandon(String target, String reason) {
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
if (waiter == null || waiter.isDone()) {
return false; // nobody is blocked on this worker — nothing to abandon
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
boolean asyncFailed = false;
for (Task task : tasks.values()) {
if (target.equals(task.target) && task.question == null
&& task.future.complete(new Reply(Outcome.WORKER_FAILED, reason))) {
asyncFailed = true;
}
}
boolean failed = rendezvous.resolveFailure(waiter, reason);
if (failed) {
log.warn("abandoning the blocked send to {}: {}", target, reason);
}
return failed;
return failed || asyncFailed;
}
/**
@@ -326,6 +359,11 @@ public final class MessageService {
* never earned. {@code null} disables the hook.
*/
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted) {
return send(target, content, timeoutMillis, onAccepted, null);
}
/** Run a send, optionally stopping an async task that teardown already failed before acceptance. */
private Reply send(String target, String content, long timeoutMillis, Runnable onAccepted, Task task) {
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
ReentrantLock lock = sessionLocks.computeIfAbsent(target, _ -> new ReentrantLock());
@@ -333,6 +371,12 @@ public final class MessageService {
return new Reply(Outcome.BUSY, null); // another send held the session the whole window
}
try {
if (task != null && task.future.isDone()) {
return task.future.getNow(null);
}
if (hasAsyncQuestion(target)) {
return new Reply(Outcome.BUSY, null); // the worker's current turn is paused for its lead
}
// Open the waiter BEFORE queueing delivery (CB-548). A fast reply — the worker already
// injectable the instant we enqueue — otherwise arrives before the waiter is registered
// and orphans into the inbox while this send blocks to the timeout (the enqueue-before-
@@ -341,13 +385,17 @@ public final class MessageService {
// failed send leaves no stale waiter behind.
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target);
try {
if (task != null) {
asyncTasksByWaiter.put(reply, task);
}
TurnToken token = new TurnToken(target, reply);
// The send has won the lock; the accepted-delivery hook records delegator ownership
// here (CB-548). It runs BEFORE enqueue so a throwing hook — onAccepted is now a
// public callback — fails the send without queuing a message that would orphan.
if (onAccepted != null) {
onAccepted.run();
}
CompletableFuture<Void> delivered = injector.enqueue(target, content);
CompletableFuture<Void> delivered = injector.enqueue(target, content, token);
try {
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId()));
@@ -364,6 +412,7 @@ public final class MessageService {
throw new IllegalStateException("interrupted awaiting reply from " + target, e);
}
} finally {
asyncTasksByWaiter.remove(reply);
rendezvous.close(target, reply);
}
} finally {
@@ -389,7 +438,12 @@ public final class MessageService {
if (ticket.fresh()) {
// Register the reverse waiter first, then surface the question — so the answer, which can
// arrive the instant the primary reacts, always finds an open waiter to resolve.
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(workerSession);
Task task = markAsyncQuestion(waiter, question, ticket.turnId());
if (!rendezvous.resolveQuestion(workerSession, question, ticket.turnId())) {
if (task != null) {
clearAsyncQuestion(ticket.turnId(), true);
}
rendezvous.closeAsk(ticket.turnId());
return new AskResult(AskOutcome.NO_WAITER, null); // no primary is blocked on this worker
}
@@ -399,6 +453,7 @@ public final class MessageService {
return new AskResult(AskOutcome.ANSWERED, answer);
} catch (TimeoutException e) {
log.debug("bridge_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
clearAsyncQuestion(ticket.turnId(), true);
return new AskResult(AskOutcome.TIMED_OUT, null);
} catch (ExecutionException e) {
Throwable cause = e.getCause();
@@ -441,9 +496,12 @@ public final class MessageService {
rendezvous.close(workerSession, reply);
return new Reply(Outcome.STALE_TURN, null); // lapsed between the lookup and the unblock
}
clearAsyncQuestion(turnId, false);
try {
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
return new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
Reply result = new Reply(outcomeOf(r.kind()), r.text(), r.turnId());
finishAsyncTask(turnId, result);
return result;
} catch (TimeoutException e) {
// The worker resumed but hasn't replied yet — no completion fallback arms an answered
// turn (it never re-entered the injector), so a silent worker rides out the window.
@@ -483,9 +541,21 @@ public final class MessageService {
*/
public String sendAsync(String target, String content, Runnable onAccepted) {
String ticket = "task-" + ticketSeq.incrementAndGet();
CompletableFuture<Reply> future = CompletableFuture.supplyAsync(
() -> send(target, content, ASYNC_TIMEOUT_MS, onAccepted), asyncExecutor);
tasks.put(ticket, new Task(target, future, System.nanoTime()));
Task task = new Task(target);
tasks.put(ticket, task);
asyncExecutor.submit(() -> {
try {
Reply result = send(target, content, ASYNC_TIMEOUT_MS, onAccepted, task);
if (result.outcome() == Outcome.QUESTION) {
// Keep the accepted owner until answer() finishes it. markAsyncQuestion may run
// just after resolveQuestion wakes this thread.
} else {
finishAsyncTask(task, result);
}
} catch (Throwable t) {
task.future.completeExceptionally(t);
}
});
pruneTerminalTickets();
log.debug("async send {} -> {}", ticket, target);
return ticket;
@@ -501,27 +571,32 @@ public final class MessageService {
if (task == null) {
return null;
}
CompletableFuture<Reply> f = task.future();
CompletableFuture<Reply> f = task.future;
if (!f.isDone()) {
return new TaskView(ticket, Phase.PENDING, null, null, "worker " + liveStatus(task.target()));
Reply question = task.question;
if (question != null) {
return new TaskView(ticket, Phase.ASKING, question.text(), null,
"worker is waiting for your answer", question.turnId());
}
return new TaskView(ticket, Phase.PENDING, null, null, "worker " + liveStatus(task.target), null);
}
Reply r;
try {
r = f.getNow(null);
} catch (CompletionException | java.util.concurrent.CancellationException e) {
Throwable cause = (e instanceof CompletionException ce && ce.getCause() != null) ? ce.getCause() : e;
return new TaskView(ticket, Phase.FAILED, null, null, cause.getMessage());
return new TaskView(ticket, Phase.FAILED, null, null, cause.getMessage(), null);
}
if (r.completed()) {
String source = r.outcome() == Outcome.REPLIED ? "reply" : "transcript";
return new TaskView(ticket, Phase.DONE, r.text(), source, null);
return new TaskView(ticket, Phase.DONE, r.text(), source, null, null);
}
// A wedged worker (CB-109) carries the error context as its reason; the timeout/busy
// outcomes carry none, so fall back to the outcome name.
String detail = r.outcome() == Outcome.WORKER_FAILED && r.text() != null
? r.text()
: "no reply — " + r.outcome().name().toLowerCase();
return new TaskView(ticket, Phase.FAILED, null, null, detail);
return new TaskView(ticket, Phase.FAILED, null, null, detail, null);
}
/** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */
@@ -536,7 +611,51 @@ public final class MessageService {
/** Drop finished tickets older than the TTL so the registry cannot grow without bound. */
private void pruneTerminalTickets() {
long cutoff = System.nanoTime() - TICKET_TTL_NANOS;
tasks.values().removeIf(t -> t.future().isDone() && t.createdNanos() < cutoff);
tasks.values().removeIf(t -> t.future.isDone() && t.createdNanos < cutoff);
}
/** Record the active question for an async ticket; blocking sends have no entry and stay unchanged. */
private Task markAsyncQuestion(CompletableFuture<Rendezvous.Resolution> waiter, String text, String turnId) {
Task task = waiter == null ? null : asyncTasksByWaiter.get(waiter);
if (task != null) {
task.question = new Reply(Outcome.QUESTION, text, turnId);
task.turnId = turnId;
asyncTasksByTurn.put(turnId, task);
}
return task;
}
/** Clear an answered or lapsed question, but only when it matches the ticket's current turn. */
private void clearAsyncQuestion(String turnId, boolean forgetTurn) {
Task task = asyncTasksByTurn.get(turnId);
if (task != null && turnId.equals(task.turnId)) {
task.question = null;
if (forgetTurn) {
asyncTasksByTurn.remove(turnId, task);
task.turnId = null;
}
}
}
/** Complete and detach an async ticket after its worker's actual terminal reply. */
private void finishAsyncTask(Task task, Reply result) {
task.future.complete(result);
if (task.turnId != null) {
asyncTasksByTurn.remove(task.turnId, task);
}
}
/** Complete the async ticket correlated to a specific answered turn. */
private void finishAsyncTask(String turnId, Reply result) {
Task task = asyncTasksByTurn.get(turnId);
if (task != null) {
finishAsyncTask(task, result);
}
}
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
private boolean hasAsyncQuestion(String target) {
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
}
/** Release the async executor. */
@@ -0,0 +1,20 @@
package dev.ltms.bridged.msg;
import java.util.concurrent.CompletableFuture;
/**
* Identity for one accepted send. The session turn is deliberately absent: CompletionResolver's
* delivery callback runs before SessionManager.onDelivered, so binding it needs a later ordering design.
*/
public final class TurnToken {
private final String target;
private final CompletableFuture<Rendezvous.Resolution> waiter;
public TurnToken(String target, CompletableFuture<Rendezvous.Resolution> waiter) {
this.target = target;
this.waiter = waiter;
}
public String target() { return target; }
public CompletableFuture<Rendezvous.Resolution> waiter() { return waiter; }
}
@@ -167,6 +167,22 @@ public final class GitWorktrees implements Worktrees {
exec("git", "-C", repoRoot, "worktree", "remove", "--force", worktreePath);
}
@Override
public boolean hasUncommitted(String worktreePath) {
// A worktree that is already gone holds no work to lose, and it must not break teardown:
// git -C <missing-dir> status exits non-zero and would throw where release() is mid-way
// through stopping a pane. Mirror remove()'s already-gone tolerance by treating it as clean.
Path p = Path.of(worktreePath);
if (!Files.exists(p)) {
log.debug("worktree {} already gone — nothing can be uncommitted", worktreePath);
return false;
}
// No --untracked-files=no: the exact shape of the work lost in CB-576 was a new file
// that was never added, so an untracked-only worktree is still dirty.
String out = exec("git", "-C", worktreePath, "status", "--porcelain");
return !out.isBlank();
}
@Override
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
if (overlay == null || overlay.isEmpty()) {
@@ -4,6 +4,7 @@ import dev.ltms.bridged.auth.MemberLifecycle;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.inject.TurnListener;
import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.msg.TurnToken;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.peer.PeerHandle;
import dev.ltms.bridged.peer.PeerLauncher;
@@ -200,6 +201,16 @@ public final class SessionManager implements TurnListener {
removed.paneId(), removed.terminalId(), removed.state(), cause);
if (preserveWorktree && removed.worktree() != null) {
logPreservedForShutdown(removed);
} else if (removed.worktree() != null && worktrees.hasUncommitted(removed.worktree())) {
// CB-576: a release that would otherwise remove the worktree finds it holding
// uncommitted work the bridge cannot see. A worker that ends a turn without
// committing (normally because it stopped to ask a question or refused the turn)
// has its only copy of that work in the worktree. Remove would --force-delete it,
// so preserve the directory and tell an operator where to find it.
preserveWorktree = true;
log.warn("release {} preserves dirty worktree {} for pane={} terminal={}: "
+ "the worktree holds uncommitted changes that --force remove would destroy",
cause, removed.worktree(), removed.paneId(), removed.terminalId());
}
// CB-516: a send still waiting on this worker can never be answered now. Tell the
// listener BEFORE the pane is torn down, so a blocked caller fails fast with a real
@@ -425,7 +436,7 @@ public final class SessionManager implements TurnListener {
* can be re-delivered for multi-turn reuse until it is released.
*/
@Override
public void onDelivered(String target) {
public void onDelivered(String target, TurnToken token) {
MemberSession current = findByTerminal(target);
if (current == null) return;
if (current.state() != MemberSession.State.READY && current.state() != MemberSession.State.DONE) {
@@ -10,6 +10,18 @@ public interface Worktrees {
/** git -C <repoRoot> worktree remove --force <path>. Idempotent (already-gone tolerated). */
void remove(String repoRoot, String worktreePath);
/**
* True when the worktree holds uncommitted changes the bridge cannot see: tracked
* modifications, staged files, or untracked files. {@code git status --porcelain} is the
* test; an empty result means clean. Callers use this to decide whether removing the
* worktree would silently destroy a worker's only copy of its work.
*
* <p>An already-gone worktree is reported as clean (no throw), matching {@link #remove}'s
* idempotent contract: a path that does not exist holds no work to lose, and must not break
* a teardown that is mid-way through stopping the pane.
*/
boolean hasUncommitted(String worktreePath);
/** Copy each existing overlay path repoRoot→worktree; mark tracked ones --skip-worktree. */
void overlayParity(String repoRoot, String worktreePath, List<String> overlay);
+9
View File
@@ -1,4 +1,13 @@
<configuration>
<!--
CB-575: MCP SDK 2.0.0 has no public notification registration API. It registers only
notifications/initialized and notifications/roots/list_changed, so other client notifications
still warn when unhandled. Clients may legitimately send notifications/cancelled; suppress only
that SDK WARN because it would devalue the action-needed WARN level used by M4 fleet health.
If a later SDK handles cancellation, this filter simply stops matching and can be removed.
-->
<turboFilter class="dev.ltms.bridged.logging.McpCancelledNotificationFilter"/>
<appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{HH:mm:ss.SSS} %-5level [%thread] %logger{28} - %msg%n</pattern>
@@ -0,0 +1,81 @@
package dev.ltms.bridged.health;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.inject.Injector;
import dev.ltms.bridged.msg.InMemoryReplyInbox;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.guard.SubscriptionGuard;
import dev.ltms.bridged.config.BridgedConfig;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.Executors;
import static org.junit.jupiter.api.Assertions.assertEquals;
class FleetHealthMonitorTest {
@Test void oneTickUsesOneFleetListForAnyRosterSize() {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
BridgedConfig.Profile profile = new BridgedConfig.Profile("test", "http://test:1", null,
null, null, null, null, null, null, null, null, null);
ClaudeCodeLauncher launcher = new ClaudeCodeLauncher(agents, new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("test")), Map.of("test", profile), "test", _ -> "token");
SessionManager sessions = new SessionManager(launcher);
sessions.acquire("test", null, null, null);
sessions.acquire("test", null, null, null);
herdr.calls.clear();
MessageService messages = new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox());
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, sessions::roster, messages, scheduler, () -> 1, 60);
monitor.tick();
monitor.stop();
assertEquals(1, herdr.calls.stream().filter(call -> call.method().equals("agent.list")).count());
}
@Test void failedTickDoesNotStopTheNextTick() {
FakeHerdr herdr = new FakeHerdr().healthy(false);
AgentControl agents = new AgentControl(herdr);
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, () -> 1, 60);
monitor.tick();
herdr.healthy(true);
monitor.tick();
monitor.stop();
assertEquals(2, herdr.calls.stream().filter(call -> call.method().equals("agent.list")).count());
}
@Test void faultTransitionLogsOnlyOnceUntilItChanges() {
Logger logger = (Logger) LoggerFactory.getLogger(FleetHealthMonitor.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
FakeHerdr herdr = new FakeHerdr();
AgentControl agents = new AgentControl(herdr);
var scheduler = Executors.newSingleThreadScheduledExecutor();
FleetHealthMonitor monitor = new FleetHealthMonitor(agents, java.util.List::of,
new MessageService(agents, new Injector(agents), new Rendezvous(), new InMemoryReplyInbox()),
scheduler, () -> 1, 60);
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
monitor.reportTransition("term_a", HealthState.TURN_BOUNDARY_LOST);
monitor.stop();
assertEquals(1, appender.list.stream().filter(event -> event.getFormattedMessage()
.contains("member=term_a state=TURN_BOUNDARY_LOST")).count());
} finally {
logger.detachAppender(appender);
}
}
}
@@ -0,0 +1,56 @@
package dev.ltms.bridged.health;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.msg.Rendezvous;
import org.junit.jupiter.api.Test;
import static org.junit.jupiter.api.Assertions.assertEquals;
class FleetHealthTest {
@Test void muteCountsOnlyCompletionFallbacks() {
MuteCounter mute = new MuteCounter();
mute.observe("target", "terra", Rendezvous.Kind.REPLY);
mute.observe("target", "terra", Rendezvous.Kind.COMPLETION);
assertEquals(1, mute.forTarget("target"));
assertEquals(1, mute.forProfile("terra"));
}
@Test void turnBoundaryNeedsTwoSnapshots() {
HealthSnapshot s = snapshot(MemberSession.State.BUSY, AgentStatus.DONE, true);
HealthDecision first = FleetHealth.decide(s, HealthPrior.NONE, 1);
assertEquals(HealthState.WORKING, first.state());
assertEquals(new HealthPrior(true), first.prior());
assertEquals(HealthState.TURN_BOUNDARY_LOST, FleetHealth.decide(s, first.prior(), 2).state());
}
@Test void unknownLiveStatusWithAcceptedDeliveryIsNotIdle() {
assertEquals(HealthState.WORKING, FleetHealth.decide(
snapshot(MemberSession.State.BUSY, AgentStatus.UNKNOWN, true), HealthPrior.NONE, 1).state());
}
@Test void acceptedDeliveryNeverReportsIdle() {
for (MemberSession.State session : MemberSession.State.values()) {
for (AgentStatus live : AgentStatus.values()) {
HealthSnapshot s = snapshot(session, live, true);
assertEquals(false, FleetHealth.decide(s, HealthPrior.NONE, 1).state() == HealthState.IDLE,
() -> "accepted delivery returned IDLE for " + session + "/" + live);
}
}
}
@Test void blockedDoesNotGuessPromptKind() {
assertEquals(HealthState.BLOCKED_AMBIGUOUS, FleetHealth.decide(
snapshot(MemberSession.State.BUSY, AgentStatus.BLOCKED, true), HealthPrior.NONE, 1).state());
}
@Test void controlLinkOutranksMemberFault() {
HealthSnapshot s = new HealthSnapshot(MemberSession.State.BUSY, AgentStatus.DONE, true, false,
false, true, true, true, false, true, true, true);
assertEquals(HealthState.CONTROL_LINK_DOWN, FleetHealth.decide(s, HealthPrior.NONE, 1).state());
}
private static HealthSnapshot snapshot(MemberSession.State state, AgentStatus live, boolean accepted) {
return new HealthSnapshot(state, live, accepted, false, false, true, false, false,
false, false, false, false);
}
}
@@ -0,0 +1,14 @@
package dev.ltms.bridged.health;
import org.junit.jupiter.api.Test;
import java.util.List;
import static org.junit.jupiter.api.Assertions.assertEquals;
class PaneBudgetTest {
@Test void fixedLimitsIgnoreWeakerConfig() {
PaneBudget budget = new PaneBudget();
List<String> targets = List.of("a", "b", "c");
assertEquals(List.of("a", "b"), budget.choose(targets, 0, 0));
assertEquals(List.of("c"), budget.choose(targets, 1, 0));
}
}
@@ -7,6 +7,8 @@ import ch.qos.logback.core.read.ListAppender;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.msg.TestTurnTokens;
import dev.ltms.bridged.msg.TurnToken;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
@@ -47,7 +49,7 @@ class CompletionResolverTest {
Rendezvous rendezvous = new Rendezvous(); // no waiter opened
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
resolver.captureBaseline("term_a"); // no send to attribute a later completion to
resolver.captureBaseline("term_a", TestTurnTokens.inert("term_a")); // no send to attribute a later completion to
assertFalse(herdr.called("agent.read"),
"with no waiting send there is no turn to baseline — skip the scrape");
@@ -194,7 +196,7 @@ class CompletionResolverTest {
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
var waiter = rendezvous.open("term_a");
resolver.captureBaseline("term_a");
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter));
herdr.readText("⏺ answer that /clear would erase\n❯ ");
resolver.resolveBeforePostAction("term_a");
@@ -217,7 +219,7 @@ class CompletionResolverTest {
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
var waiter = rendezvous.open("term_a"); // a send is blocked on this turn
resolver.captureBaseline("term_a"); // baseline is the clipped >cap block
resolver.captureBaseline("term_a", new TurnToken("term_a", waiter)); // baseline is the clipped >cap block
var turn = resolver.inFlight("term_a");
assertEquals(CompletionResolver.MAX_SCRAPE_CHARS, turn.baseline().length(),
"the delivery baseline is clipped to the same cap resolve() applies to the tail");
@@ -8,6 +8,7 @@ import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.msg.TestTurnTokens;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
@@ -47,7 +48,7 @@ class InjectorTest {
@Test
void deliversWhenIdle() {
CompletableFuture<Void> f = injector.enqueue(T, "hello");
CompletableFuture<Void> f = injector.enqueue(T, "hello", TestTurnTokens.inert(T));
assertFalse(f.isDone(), "not delivered until an injectable status arrives");
injector.onStatus(T, AgentStatus.IDLE);
assertTrue(f.isDone());
@@ -59,7 +60,7 @@ class InjectorTest {
// CB-113: idle alone is not enough — hold until the worker's MCP is connected (ready).
java.util.Set<String> ready = new java.util.HashSet<>();
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, ready::contains);
inj.enqueue(T, "task");
inj.enqueue(T, "task", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // idle but not yet available → held out of the boot window
assertEquals(List.of(), sent(), "must not deliver into a not-yet-available worker");
@@ -79,7 +80,7 @@ class InjectorTest {
void resubmitsEnterWhenADeliveredMessageIsNotPickedUp() {
// CB-113: the Enter at delivery can race the paste; while the worker stays idle (not picked
// up), the injector re-nudges Enter so the pending paste submits.
injector.enqueue(T, "task");
injector.enqueue(T, "task", TestTurnTokens.inert(T));
injector.onStatus(T, AgentStatus.IDLE); // deliver: paste + one Enter
long afterDeliver = enterKeystrokes();
@@ -95,7 +96,7 @@ class InjectorTest {
@Test
void holdsWhileWorkingThenDeliversOnIdle() {
injector.enqueue(T, "later");
injector.enqueue(T, "later", TestTurnTokens.inert(T));
injector.onStatus(T, AgentStatus.WORKING);
assertEquals(List.of(), sent(), "must not inject mid-turn");
injector.onStatus(T, AgentStatus.IDLE);
@@ -104,7 +105,7 @@ class InjectorTest {
@Test
void blockedIsInjectableButUnknownIsNot() {
injector.enqueue(T, "answer");
injector.enqueue(T, "answer", TestTurnTokens.inert(T));
injector.onStatus(T, AgentStatus.UNKNOWN);
assertEquals(List.of(), sent(), "unknown status is not safe to inject");
injector.onStatus(T, AgentStatus.BLOCKED);
@@ -113,8 +114,8 @@ class InjectorTest {
@Test
void twoRapidDeliveriesNeverInterleave() {
injector.enqueue(T, "m1");
injector.enqueue(T, "m2");
injector.enqueue(T, "m1", TestTurnTokens.inert(T));
injector.enqueue(T, "m2", TestTurnTokens.inert(T));
// First idle window delivers only m1, even if idle is observed twice before pickup.
injector.onStatus(T, AgentStatus.IDLE);
@@ -129,8 +130,8 @@ class InjectorTest {
@Test
void transientUnknownDoesNotReleaseThePickupLatch() {
injector.enqueue(T, "m1");
injector.enqueue(T, "m2");
injector.enqueue(T, "m1", TestTurnTokens.inert(T));
injector.enqueue(T, "m2", TestTurnTokens.inert(T));
injector.onStatus(T, AgentStatus.IDLE); // m1 sent, awaiting pickup
assertEquals(List.of("m1"), sent());
@@ -145,8 +146,8 @@ class InjectorTest {
@Test
void missedPickupEdgeIsReleasedByGraceSoTheQueueNeverWedges() {
injector.enqueue(T, "m1");
injector.enqueue(T, "m2");
injector.enqueue(T, "m1", TestTurnTokens.inert(T));
injector.enqueue(T, "m2", TestTurnTokens.inert(T));
injector.onStatus(T, AgentStatus.IDLE); // m1 sent
assertEquals(List.of("m1"), sent());
@@ -158,9 +159,9 @@ class InjectorTest {
@Test
void fifoOrderAcrossManyTurns() {
injector.enqueue(T, "a");
injector.enqueue(T, "b");
injector.enqueue(T, "c");
injector.enqueue(T, "a", TestTurnTokens.inert(T));
injector.enqueue(T, "b", TestTurnTokens.inert(T));
injector.enqueue(T, "c", TestTurnTokens.inert(T));
for (int i = 0; i < 3; i++) {
injector.onStatus(T, AgentStatus.IDLE); // deliver one
injector.onStatus(T, AgentStatus.WORKING); // pickup
@@ -172,7 +173,7 @@ class InjectorTest {
@Test
void activeWhileQueuedOrInFlightThenQuietAfterTurnCompletes() {
assertTrue(injector.activeTargets().isEmpty());
injector.enqueue(T, "x");
injector.enqueue(T, "x", TestTurnTokens.inert(T));
assertEquals(Set.of(T), injector.activeTargets(), "active while a message is queued");
injector.onStatus(T, AgentStatus.IDLE); // delivers; awaiting pickup
@@ -191,7 +192,7 @@ class InjectorTest {
void firesTurnCompleteOnAConfirmedWorkingThenIdle() {
List<String> completed = new ArrayList<>();
Injector inj = new Injector(new AgentControl(herdr), completed::add);
inj.enqueue(T, "task");
inj.enqueue(T, "task", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // deliver
inj.onStatus(T, AgentStatus.WORKING); // pickup + turn running
@@ -226,8 +227,8 @@ class InjectorTest {
}
ResetListener listener = new ResetListener();
Injector inj = new Injector(agents, listener);
inj.enqueue(T, "first");
inj.enqueue(T, "second");
inj.enqueue(T, "first", TestTurnTokens.inert(T));
inj.enqueue(T, "second", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // first delegation
inj.onStatus(T, AgentStatus.WORKING);
@@ -246,7 +247,7 @@ class InjectorTest {
void doesNotSynthesizeCompletionFromAnUnconfirmedTurn() {
List<String> completed = new ArrayList<>();
Injector inj = new Injector(new AgentControl(herdr), completed::add);
inj.enqueue(T, "task");
inj.enqueue(T, "task", TestTurnTokens.inert(T));
// Deliver, then only ever idle — a `working` sample is never seen. The pickup grace unwedges
// the queue but must NOT invent a completion: without a sampled turn there is no trustworthy
@@ -259,6 +260,7 @@ class InjectorTest {
private static final class Captor implements TurnListener {
final List<String> completed = new ArrayList<>();
final List<String> failed = new ArrayList<>();
final List<String> failureReasons = new ArrayList<>();
@Override
public void onTurnComplete(String target) {
@@ -269,6 +271,12 @@ class InjectorTest {
public void onTurnFailed(String target) {
failed.add(target);
}
@Override
public void onTurnFailed(String target, String reason) {
failed.add(target);
failureReasons.add(reason);
}
}
// ~30s of unknown at the 250ms prod poll interval; enough onStatus samples to trip the stall.
@@ -278,7 +286,7 @@ class InjectorTest {
void failsAnOutstandingDelegationWhoseWorkerWedgesInUnknown() {
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), cap);
inj.enqueue(T, "task");
inj.enqueue(T, "task", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // deliver
inj.onStatus(T, AgentStatus.WORKING); // worker starts the turn
@@ -293,7 +301,7 @@ class InjectorTest {
void aTransientUnknownGlitchNeitherFailsNorBlocksCompletion() {
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), cap);
inj.enqueue(T, "task");
inj.enqueue(T, "task", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // deliver
inj.onStatus(T, AgentStatus.WORKING); // confirmed turn
@@ -308,7 +316,7 @@ class InjectorTest {
void sendFailureDropsMessageAndFailsItsFuture() {
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
Injector inj = new Injector(new AgentControl(failing));
CompletableFuture<Void> f = inj.enqueue(T, "boom");
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE);
assertTrue(f.isCompletedExceptionally());
@@ -317,18 +325,35 @@ class InjectorTest {
@Test
void dropFailsPendingWaiters() {
CompletableFuture<Void> f = injector.enqueue(T, "orphan");
CompletableFuture<Void> f = injector.enqueue(T, "orphan", TestTurnTokens.inert(T));
injector.drop(T, new HerdrException("worker gone", "pane_not_found", null));
assertTrue(f.isCompletedExceptionally(), "queued waiters unblock when the worker vanishes");
}
@Test
void dropPassesTheRealCauseForQueuedAndDeliveredWork() {
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), cap);
CompletableFuture<Void> delivered = inj.enqueue(T, "delivered", TestTurnTokens.inert(T));
CompletableFuture<Void> queued = inj.enqueue(T, "queued", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // deliver the first message
inj.onStatus(T, AgentStatus.WORKING); // its turn is now in flight; one remains queued
inj.drop(T, new HerdrException("agent target sol not found", "agent_not_found", null));
assertEquals(List.of(T), cap.failed, "drop signals one turn failure for both affected states");
assertEquals(List.of("agent target sol not found"), cap.failureReasons);
assertTrue(delivered.isDone(), "the delivered future has already completed");
assertTrue(queued.isCompletedExceptionally(), "the queued future fails with the drop cause");
}
@Test
void dropFailsTheTurnOfADeliveredMessageWhenTheWorkerVanishes() {
// CB-110: the message was delivered (no longer queued), so failing queued waiters alone would
// leave its send hanging. A vanished worker must fail that in-flight turn too.
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), cap);
inj.enqueue(T, "task");
inj.enqueue(T, "task", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // deliver
inj.onStatus(T, AgentStatus.WORKING); // turn running
@@ -343,7 +368,7 @@ class InjectorTest {
// by an off-sub worker's review of CB-110, delegated through the bridge.)
Captor cap = new Captor();
Injector inj = new Injector(new AgentControl(herdr), cap);
inj.enqueue(T, "task");
inj.enqueue(T, "task", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE); // deliver; pickup never confirmed
inj.drop(T, new HerdrException("worker gone", "pane_not_found", null));
@@ -362,7 +387,7 @@ class InjectorTest {
Captor cap = new Captor();
List<String> forgotten = new ArrayList<>();
Injector inj = new Injector(new AgentControl(herdr), cap, _ -> false, forgotten::add);
CompletableFuture<Void> f = inj.enqueue(T, "task");
CompletableFuture<Void> f = inj.enqueue(T, "task", TestTurnTokens.inert(T));
for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE);
@@ -390,7 +415,7 @@ class InjectorTest {
try {
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> false, _ -> {
});
inj.enqueue(T, "task");
inj.enqueue(T, "task", TestTurnTokens.inert(T));
for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE);
@@ -415,7 +440,7 @@ class InjectorTest {
Set<String> ready = new java.util.HashSet<>();
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, ready::contains, _ -> {
});
inj.enqueue(T, "task");
inj.enqueue(T, "task", TestTurnTokens.inert(T));
for (int i = 0; i < 100; i++) inj.onStatus(T, AgentStatus.IDLE); // still booting, well under grace
assertEquals(List.of(), sent());
@@ -431,7 +456,7 @@ class InjectorTest {
// linger past the worker's life (MemberPresence.forget had no caller before this).
List<String> forgotten = new ArrayList<>();
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> true, forgotten::add);
inj.enqueue(T, "orphan");
inj.enqueue(T, "orphan", TestTurnTokens.inert(T));
inj.drop(T, new HerdrException("worker gone", "pane_not_found", null));
assertEquals(List.of(T), forgotten, "drop clears the gone worker's presence");
}
@@ -452,7 +477,7 @@ class InjectorTest {
try {
Injector inj = new Injector(new AgentControl(herdr), TurnListener.NOOP, _ -> true, _ -> {
});
inj.enqueue(T, "orphan");
inj.enqueue(T, "orphan", TestTurnTokens.inert(T));
inj.drop(T, new HerdrException("worker gone", "pane_not_found", null));
String warn = appender.list.stream()
@@ -476,7 +501,7 @@ class InjectorTest {
StatusPoller poller = new StatusPoller(new AgentControl(idle), inj, 10);
poller.start();
try {
CompletableFuture<Void> delivered = inj.enqueue(T, "via-poller");
CompletableFuture<Void> delivered = inj.enqueue(T, "via-poller", TestTurnTokens.inert(T));
delivered.get(2, TimeUnit.SECONDS); // completes when the poller drives the send
} finally {
poller.stop();
@@ -495,7 +520,7 @@ class InjectorTest {
void deliveredFutureCarriesSendFailure() {
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
Injector inj = new Injector(new AgentControl(failing));
CompletableFuture<Void> f = inj.enqueue(T, "boom");
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
inj.onStatus(T, AgentStatus.IDLE);
ExecutionException ex = assertThrows(ExecutionException.class, f::get);
assertInstanceOf(HerdrException.class, ex.getCause());
@@ -0,0 +1,37 @@
package dev.ltms.bridged.logging;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import io.modelcontextprotocol.spec.McpSchema;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
class McpCancelledNotificationFilterTest {
@Test
void suppressesOnlyTheCancelledNotificationWarning() {
Logger logger = (Logger) LoggerFactory.getLogger(McpCancelledNotificationFilter.LOGGER);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.start();
logger.addAppender(appender);
try {
logger.warn(McpCancelledNotificationFilter.UNHANDLED_NOTIFICATION,
new McpSchema.JSONRPCNotification("notifications/cancelled", Map.of("requestId", 7)));
logger.warn(McpCancelledNotificationFilter.UNHANDLED_NOTIFICATION,
new McpSchema.JSONRPCNotification("notifications/progress", Map.of("progress", 1)));
assertEquals(1, appender.list.size());
assertEquals(Level.WARN, appender.list.getFirst().getLevel());
assertTrue(appender.list.getFirst().getFormattedMessage().contains("notifications/progress"));
} finally {
logger.detachAppender(appender);
}
}
}
@@ -72,7 +72,7 @@ class BridgeMcpAuthzTest {
new PrimaryRegistry(null),
enforce ? CallerResolver.withLeadsAndMembers(identity, false, null,
Map::of, new MemberRegistry(null)) : null,
metrics);
metrics, BridgeMcp.CapacitySource.none(), new BridgeMcp.HealthCoverageSource(() -> "off"));
return mcp;
}
@@ -53,6 +53,18 @@ class BridgeMcpTest {
return ((McpSchema.TextContent) r.content().getFirst()).text();
}
private void assertSendRoundTrips(String target, Set<String> profiles) throws Exception {
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> BridgeMcp.send(messages, target, "hi", 4000L, null, profiles));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting(target) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting(target), "send should be accepted for " + target);
BridgeMcp.reply(messages, target, "received");
assertEquals("received", textOf(send.get(6, TimeUnit.SECONDS)));
}
private static ClaudeCodeLauncher workerService(FakeHerdr h, String baseUrl, Set<String> allow) {
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
"ltms-local", baseUrl, "coder", null, "BRIDGED_WORKER_TOKEN", null,
@@ -69,7 +81,7 @@ class BridgeMcpTest {
void sendThenReplyRoundTrips() throws Exception {
// bridge_send blocks; bridge_reply resolves it with the worker's structured answer.
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> BridgeMcp.send(messages, "term_a", "review this", 4000L));
() -> BridgeMcp.send(messages, "term_a", "review this", 4000L, null, Set.of()));
// Wait until the send has opened its waiter so the reply resolves it (CB-307: reply now
// queues in the inbox if no waiter is open, which would break the round-trip).
@@ -91,7 +103,7 @@ class BridgeMcpTest {
@Test
void asyncSendReturnsATicketThenPollReportsTheReply() throws Exception {
// wait:false parity — a ticket is issued, resolved by a reply, and surfaced by bridge_poll.
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it");
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
assertNotEquals(Boolean.TRUE, accepted.isError());
String out = textOf(accepted);
assertTrue(out.contains("ticket="), out);
@@ -120,6 +132,131 @@ class BridgeMcpTest {
assertEquals("async LGTM", textOf(polled));
}
@Test
void asyncSendSurfacesAnAskThenKeepsTheTicketForTheFinalReply() throws Exception {
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"));
CompletableFuture<McpSchema.CallToolResult> ask = CompletableFuture.supplyAsync(
() -> BridgeMcp.ask(messages, "term_a", "which config?", 5000L));
McpSchema.CallToolResult question = BridgeMcp.poll(messages, ticket, null);
deadline = System.currentTimeMillis() + 3000;
while (!textOf(question).contains("[question") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
question = BridgeMcp.poll(messages, ticket, null);
}
assertTrue(textOf(question).contains("which config?"), textOf(question));
String questionText = textOf(question);
String afterTurnId = questionText.substring(questionText.indexOf("turnId=\"") + "turnId=\"".length());
String turnId = afterTurnId.substring(0, afterTurnId.indexOf('"'));
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> BridgeMcp.answer(messages, turnId, "config.yaml", 5000L));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"));
BridgeMcp.reply(messages, "term_a", "done");
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
McpSchema.CallToolResult done = BridgeMcp.poll(messages, ticket, null);
deadline = System.currentTimeMillis() + 3000;
while (!"done".equals(textOf(done)) && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
done = BridgeMcp.poll(messages, ticket, null);
}
assertEquals("done", textOf(done));
}
@Test
void unansweredAsyncAskReturnsTheTicketToPending() throws Exception {
McpSchema.CallToolResult accepted = BridgeMcp.sendAsync(messages, "term_a", "do it", null, Set.of());
String ticket = textOf(accepted).substring(textOf(accepted).indexOf("ticket=") + "ticket=".length()).trim();
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"));
McpSchema.CallToolResult ask = BridgeMcp.ask(messages, "term_a", "still there?", 50L);
assertTrue(textOf(ask).contains("no answer"), textOf(ask));
assertTrue(textOf(BridgeMcp.poll(messages, ticket, null)).startsWith("[pending"));
BridgeMcp.reply(messages, "term_a", "finished after timeout");
assertEquals("finished after timeout", messages.drainReplies("term_a").getFirst().content());
}
@Test
void asyncSendFailureDoesNotLeaveItsTicketPending() throws Exception {
String ticket = messages.sendAsync("term_a", "do it", () -> {
throw new IllegalStateException("accept failed");
});
long deadline = System.currentTimeMillis() + 3000;
MessageService.TaskView view = messages.poll(ticket);
while (view.phase() == MessageService.Phase.PENDING && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
view = messages.poll(ticket);
}
assertEquals(MessageService.Phase.FAILED, view.phase());
assertEquals("accept failed", view.detail());
}
@Test
void anotherAsyncTicketCannotCaptureAReplyWhileTheFirstTicketIsAsking() throws Exception {
McpSchema.CallToolResult firstAccepted = BridgeMcp.sendAsync(messages, "term_a", "first", null, Set.of());
String firstTicket = textOf(firstAccepted).substring(textOf(firstAccepted).indexOf("ticket=") + "ticket=".length()).trim();
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
assertTrue(rendezvous.isWaiting("term_a"));
CompletableFuture<McpSchema.CallToolResult> ask = CompletableFuture.supplyAsync(
() -> BridgeMcp.ask(messages, "term_a", "which config?", 5000L));
MessageService.TaskView first = messages.poll(firstTicket);
deadline = System.currentTimeMillis() + 3000;
while (first.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
first = messages.poll(firstTicket);
}
assertEquals(MessageService.Phase.ASKING, first.phase());
String firstTurnId = first.turnId();
McpSchema.CallToolResult secondAccepted = BridgeMcp.sendAsync(messages, "term_a", "second", null, Set.of());
String secondTicket = textOf(secondAccepted).substring(textOf(secondAccepted).indexOf("ticket=") + "ticket=".length()).trim();
MessageService.TaskView second = messages.poll(secondTicket);
deadline = System.currentTimeMillis() + 3000;
while (second.phase() == MessageService.Phase.PENDING && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
second = messages.poll(secondTicket);
}
assertEquals(MessageService.Phase.FAILED, second.phase());
BridgeMcp.reply(messages, "term_a", "late reply");
assertEquals("late reply", messages.drainReplies("term_a").getFirst().content());
CompletableFuture<McpSchema.CallToolResult> answer = CompletableFuture.supplyAsync(
() -> BridgeMcp.answer(messages, firstTurnId, "config.yaml", 5000L));
assertEquals("config.yaml", textOf(ask.get(6, TimeUnit.SECONDS)));
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
Thread.sleep(5);
}
BridgeMcp.reply(messages, "term_a", "done");
assertEquals("done", textOf(answer.get(6, TimeUnit.SECONDS)));
}
@Test
void pollUnknownTicketIsAnError() {
McpSchema.CallToolResult res = BridgeMcp.poll(messages, "task-999", null);
@@ -129,15 +266,40 @@ class BridgeMcpTest {
@Test
void sendTimesOutWithAWorkingNote() {
McpSchema.CallToolResult res = BridgeMcp.send(messages, "term_a", "hi", 120L);
McpSchema.CallToolResult res = BridgeMcp.send(messages, "term_a", "hi", 120L, null, Set.of());
assertNotEquals(Boolean.TRUE, res.isError(), "a timeout is informational, not a tool error");
assertTrue(textOf(res).contains("no reply"), "got: " + textOf(res));
}
@Test
void sendRejectsMissingArgs() {
assertTrue(BridgeMcp.send(messages, null, "hi", null).isError());
assertTrue(BridgeMcp.send(messages, "term_a", " ", null).isError());
assertTrue(BridgeMcp.send(messages, null, "hi", null, null, Set.of()).isError());
assertTrue(BridgeMcp.send(messages, "term_a", " ", null, null, Set.of()).isError());
}
@Test
void sendRejectsAConfiguredProfileNameBeforeAcceptingIt() {
McpSchema.CallToolResult blocking = BridgeMcp.send(messages, "sol", "hi", 100L, null, Set.of("sol"));
McpSchema.CallToolResult async = BridgeMcp.sendAsync(messages, "sol", "hi", null, Set.of("sol"));
assertTrue(blocking.isError());
assertTrue(async.isError());
assertTrue(textOf(blocking).contains("sol"));
assertTrue(textOf(blocking).contains("configured profile name"));
assertTrue(textOf(blocking).contains("bridge_list"));
assertFalse(textOf(async).contains("ticket="));
}
@Test
void sendAllowsPeerLeadMemberAndUnclassifiedTargets() throws Exception {
Set<String> profiles = Set.of("sol");
assertSendRoundTrips("term_peer_lead", profiles);
assertSendRoundTrips("term_live_member", profiles);
// A herdr-owned pane outside the bridge roster cannot be classified at accept time.
McpSchema.CallToolResult result = BridgeMcp.send(messages, "external-pane", "hi", 10L, null, profiles);
assertFalse(result.isError(), "an unclassified target must not be rejected at acceptance time");
}
@Test
@@ -173,7 +335,7 @@ class BridgeMcpTest {
void askThenAnswerRoundTrips() throws Exception {
// The primary delegates and blocks; wait until its waiter is open before the worker asks.
CompletableFuture<McpSchema.CallToolResult> send = CompletableFuture.supplyAsync(
() -> BridgeMcp.send(messages, "term_a", "do X", 5000L));
() -> BridgeMcp.send(messages, "term_a", "do X", 5000L, null, Set.of()));
long deadline = System.currentTimeMillis() + 3000;
while (!rendezvous.isWaiting("term_a") && System.currentTimeMillis() < deadline) {
//noinspection BusyWait
@@ -305,6 +467,42 @@ class BridgeMcpTest {
assertTrue(out.contains("\"liveStatus\":\"unknown\""), out);
}
@Test
void capacityUsesThePlacementLiveCount() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
sessions.acquire("ltms-local", null, null, null);
McpSchema.CallToolResult res = BridgeMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
sessions, null, new BridgeMcp.CapacitySource(profile -> 2, profile -> 2,
() -> Set.of("ltms-local"), () -> 0), new BridgeMcp.HealthCoverageSource(() -> "off"), Map.of(), "");
String out = textOf(res);
assertTrue(out.contains("\"maxLoad\":2"), out);
assertTrue(out.contains("\"live\":2"), out);
assertTrue(out.contains("\"free\":0"), out);
}
@Test
void capacityIncludesConfiguredProfileWithoutMembers() {
FakeHerdr h = new FakeHerdr();
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
String out = textOf(BridgeMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
sessions, null, new BridgeMcp.CapacitySource(profile -> 0, profile -> 2,
() -> Set.of("terra"), () -> 0), new BridgeMcp.HealthCoverageSource(() -> "off"), Map.of(), ""));
assertTrue(out.contains("\"profile\":\"terra\""), out);
assertTrue(out.contains("\"live\":0"), out);
assertTrue(out.contains("\"free\":2"), out);
assertTrue(out.contains("\"reclaimable\":0"), out);
}
@Test
void inertCapacitySourceOmitsCapacityBlock() {
FakeHerdr h = new FakeHerdr();
String out = textOf(BridgeMcp.listFleet(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")),
new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))), null,
BridgeMcp.CapacitySource.none(), new BridgeMcp.HealthCoverageSource(() -> "off"), Map.of(), ""));
assertFalse(out.contains("\"capacity\":"), out);
}
@Test
void listReportsLeadsAndFlagsTheCallersOwnRow() {
FakeHerdr h = new FakeHerdr();
@@ -578,7 +776,7 @@ class BridgeMcpTest {
assertTrue(out.contains("\"role\":\"dev\""), out);
assertTrue(out.contains("\"role\":\"reviewer\""), out);
assertEquals(2, out.split("\"profile\":\"ltms-local\"", -1).length - 1,
"both members share one profile — that is the point: " + out);
"inert capacity is omitted, leaving the two member rows: " + out);
}
@Test
@@ -102,6 +102,53 @@ class ClaudeCodeLauncherTest {
"the operator's own args are preserved, in order, ahead of the model flag");
}
@Test
void appendsTheBaseComposedRoleAndReplyCharter() {
FakeHerdr herdr = new FakeHerdr();
String roleCharter = "You review changes.";
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
"sonnet", "http://gx00.gw:8000", null, null, "BRIDGED_WORKER_TOKEN",
List.of("claude"), "tab", "bridged-workers", "w #{n}",
"http://127.0.0.1:8765/mcp", null, null);
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null,
0, 0L, () -> fleet(Map.of("reviewer", roleCharter), null));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.REVIEWER));
List<String> args = spawnedArgs(herdr);
int flag = args.indexOf("--append-system-prompt");
assertEquals(1, args.stream().filter("--append-system-prompt"::equals).count(),
"the composed charter is passed once");
assertEquals(roleCharter + "\n\n" + HerdrPeerLauncher.REPLY_CHARTER, args.get(flag + 1),
"the role charter comes first and the reply rule comes last");
}
@Test
void profileWithoutMcpOrRoleCharterGetsNoSystemPrompt() {
FakeHerdr herdr = new FakeHerdr();
ClaudeCodeLauncher svc = labelService(herdr, () -> fleet(Map.of(), null));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.DEV));
assertFalse(spawnedArgs(herdr).contains("--append-system-prompt"));
}
@Test
void profileWithoutMcpStillGetsItsRoleCharter() {
FakeHerdr herdr = new FakeHerdr();
String roleCharter = "You design changes.";
ClaudeCodeLauncher svc = labelService(herdr, () -> fleet(Map.of("architect", roleCharter), null));
svc.spawn(new SpawnRequest("sonnet", null, null, null, null, MemberRole.ARCHITECT));
List<String> args = spawnedArgs(herdr);
int flag = args.indexOf("--append-system-prompt");
assertTrue(flag >= 0, "a role charter does not need an MCP mount");
assertEquals(roleCharter, args.get(flag + 1));
assertFalse(args.contains("--mcp-config"));
}
private ClaudeCodeLauncher multiProfile(FakeHerdr herdr) {
BridgedConfig.Profile gx10 = new BridgedConfig.Profile("gx10", "http://gx10.gw:8000", "coder",
null, "BRIDGED_WORKER_TOKEN", List.of("claude"), "tab", "bridged-workers", "w #{n}", null, null, null);
@@ -120,6 +120,36 @@ class MessageServiceTest {
assertFalse(reply.completed());
}
@Test
void droppedQueuedAndDeliveredTurnsExposeTheRealCauseExactlyOnce() throws Exception {
CompletableFuture<MessageService.Reply> first = sendAsync();
awaitWaiting();
injector.onStatus(T, AgentStatus.IDLE); // first delivery
injector.onStatus(T, AgentStatus.WORKING); // first turn in flight
CompletableFuture<Void> queued = injector.enqueue(T, "second task", TestTurnTokens.inert(T));
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(T);
injector.drop(T, new HerdrException("agent target sol not found", "agent_not_found", null));
MessageService.Reply reply = first.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.WORKER_FAILED, reply.outcome());
assertEquals("agent target sol not found", reply.text());
assertTrue(queued.isCompletedExceptionally(), "the queued delivery future also fails");
assertFalse(rendezvous.resolveFailure(waiter, "second failure"), "the waiter fails exactly once");
}
@Test
void noDropReasonKeepsTheExistingFallbackText() throws Exception {
herdr.readText("");
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open(T);
completion.onTurnFailed(T);
Rendezvous.Resolution resolution = waiter.get(2, TimeUnit.SECONDS);
assertEquals("worker did not reply; its turn ended in an unrecoverable state "
+ "(worker unreachable or stuck)", resolution.text());
}
// --- bridge_ask reverse rendezvous (CB-205) ------------------------------------------------
@Test
@@ -609,4 +639,94 @@ class MessageServiceTest {
assertTrue(view.detail() != null && view.detail().contains("released"),
"and the detail says why, rather than 'worker unknown'");
}
@Test
void abandonFailsEveryPendingAsyncTicketForTheReleasedTarget() throws Exception {
String first = messages.sendAsync(T, "first task");
awaitWaiting(); // first task owns the target lock and rendezvous waiter
String second = messages.sendAsync(T, "second task"); // parked on the same lock, not yet queued
String third = messages.sendAsync(T, "third task"); // a second queued ticket proves the full sweep
assertTrue(messages.abandon(T, "agent target term_a not found"));
assertFailedTicket(first, "agent target term_a not found");
assertFailedTicket(second, "agent target term_a not found");
assertFailedTicket(third, "agent target term_a not found");
}
@Test
void abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
assertFalse(messages.abandon(T, "agent target term_a not found"),
"an asking ticket is an active turn, not a pending send to sweep");
assertEquals(MessageService.Phase.ASKING, messages.poll(ticket).phase());
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
}
@Test
void unansweredAsyncQuestionReturnsTheTicketToPendingAndReleasesItsTarget() throws Exception {
String ticket = messages.sendAsync(T, "task that asks");
awaitWaiting();
injectDelivery();
assertEquals(MessageService.AskOutcome.TIMED_OUT,
messages.ask(T, "which config?", 200).outcome());
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(),
"only the question wait ended; the delegated turn may still finish");
String next = messages.sendAsync(T, "next task");
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
assertEquals(MessageService.Phase.DONE, awaitTicketPhase(next, MessageService.Phase.DONE).phase());
}
@Test
void asyncQuestionBelongsToTheTaskThatOwnsItsForwardWaiter() throws Exception {
String first = messages.sendAsync(T, "first task");
awaitWaiting();
CompletableFuture<MessageService.AskResult> ask =
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
assertEquals(MessageService.Phase.ASKING, awaitTicketPhase(first, MessageService.Phase.ASKING).phase());
MessageService.TaskView asking = messages.poll(first);
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
awaitWaiting();
assertTrue(rendezvous.resolve(T, "done"));
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
}
private void assertFailedTicket(String ticket, String reason) throws Exception {
MessageService.TaskView view = awaitTicketPhase(ticket, MessageService.Phase.FAILED);
assertEquals(reason, view.detail());
}
private MessageService.TaskView awaitTicketPhase(String ticket, MessageService.Phase phase) throws Exception {
long deadline = System.currentTimeMillis() + 3000;
MessageService.TaskView view;
do {
view = messages.poll(ticket);
if (view.phase() == phase) {
return view;
}
Thread.sleep(5);
} while (System.currentTimeMillis() < deadline);
assertEquals(phase, view.phase());
return view;
}
}
@@ -0,0 +1,20 @@
package dev.ltms.bridged.msg;
/**
* Explicit unbound tokens for tests that exercise delivery without an accepted send.
*
* <p>The waiter is {@code null} on purpose. "No accepted send" is an <em>absence</em>, and a helper
* that handed back a fresh {@code CompletableFuture} would invent one — which is how the first
* version of this class turned {@code captureBaselineSkipsTheReadWhenNoSendIsWaiting} red: the
* resolver saw a non-null waiter, decided a turn was in flight, and scraped a pane that no send was
* blocked on. An inert value must omit the fact, never fabricate it.
*/
public final class TestTurnTokens {
private TestTurnTokens() {
}
/** A token for a delivery that no send is waiting on: it authorises nothing. */
public static TurnToken inert(String target) {
return new TurnToken(target, null);
}
}
@@ -29,8 +29,11 @@ public final class FakeWorktrees implements Worktrees {
private final Set<String> existingPaths = ConcurrentHashMap.newKeySet();
private final Set<String> trackedPaths = ConcurrentHashMap.newKeySet();
private volatile RuntimeException addFailure;
private volatile boolean dirty = false;
private volatile String repoRoot = "/repo";
private volatile String prefix = "/worktrees";
/** Worktree paths that currently exist, mirroring real {@code Files.exists} for the gone case. */
private final Set<String> worktreePaths = ConcurrentHashMap.newKeySet();
public FakeWorktrees withRepoRoot(String root) {
this.repoRoot = root;
@@ -61,6 +64,12 @@ public final class FakeWorktrees implements Worktrees {
return this;
}
/** Mark the worktree dirty so {@link #hasUncommitted} reports true (simulates uncommitted work). */
public FakeWorktrees withDirty(boolean dirty) {
this.dirty = dirty;
return this;
}
@Override
public String add(String repoRoot, String branch, String baseRef) {
addCalls.add(new AddCall(repoRoot, branch, baseRef));
@@ -69,7 +78,15 @@ public final class FakeWorktrees implements Worktrees {
}
// The branch already carries a unique nonce, so the derived path is distinct per acquire
// without an extra counter — keep it a pure function of the branch the test can predict.
return prefix + "/" + branch.replace('/', '_');
String path = prefix + "/" + branch.replace('/', '_');
worktreePaths.add(path);
return path;
}
/** Model an operator / {@code git worktree prune} removing the worktree before release. */
public FakeWorktrees markGone(String worktreePath) {
worktreePaths.remove(worktreePath);
return this;
}
@Override
@@ -77,6 +94,16 @@ public final class FakeWorktrees implements Worktrees {
removeCalls.add(new RemoveCall(repoRoot, worktreePath));
}
@Override
public boolean hasUncommitted(String worktreePath) {
// A path that does not exist (never added, or marked gone) is reported clean, mirroring
// GitWorktrees' already-gone guard — never an error, so teardown still completes.
if (!worktreePaths.contains(worktreePath)) {
return false;
}
return dirty;
}
@Override
public void overlayParity(String repoRoot, String worktreePath, List<String> overlay) {
List<String> copied = new java.util.ArrayList<>();
@@ -175,6 +175,46 @@ class GitWorktreesTest {
assertTrue(Files.exists(Path.of(wt).resolve(".mcp.json")), ".mcp.json stub was dropped");
}
/**
* CB-576. {@code hasUncommitted} must treat a freshly-provisioned worktree as clean, but a
* worktree holding a brand-new, never-added file as dirty. The untracked-file-only shape is
* exactly the work lost in the incident — a worker's draft that compiled but was never
* committed because it stopped to ask its lead a question.
*/
@Test
void anUntrackedOnlyWorktreeCountsAsDirty(@TempDir Path tmp) throws Exception {
Path repo = initRepo(tmp.resolve("repo"));
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
String wt = gitWorktrees.add(repo.toString(), "cb-576-u", "HEAD");
assertFalse(gitWorktrees.hasUncommitted(wt),
"a freshly provisioned worktree must read as clean");
Files.writeString(Path.of(wt).resolve("brand-new.txt"), "draft that was never added\n");
assertTrue(gitWorktrees.hasUncommitted(wt),
"an untracked-only file must count as dirty");
Files.writeString(Path.of(wt).resolve("README.md"), "edited tracked file\n");
assertTrue(gitWorktrees.hasUncommitted(wt),
"a tracked modification must also count as dirty");
}
/**
* CB-576 review. {@code hasUncommitted} must tolerate a missing worktree exactly like
* {@code remove}: an already-gone directory holds no work to lose, and throwing here would
* break teardown — SessionManager.release() calls it before stopping the pane, so an
* exception would orphan a live pane and skip the release notification (CB-516).
*/
@Test
void hasUncommittedOnAMissingWorktreeReturnsFalseWithoutThrowing(@TempDir Path tmp) {
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
String gone = tmp.resolve("wts").resolve("does-not-exist").toString();
assertFalse(gitWorktrees.hasUncommitted(gone),
"a missing worktree is reported clean, not an error");
}
/** All three protected configs are covered: each one present in a worktree is neutralized and hidden. */
@Test
void allThreeConfigsAreNeutralizedWhenPresent(@TempDir Path tmp) throws Exception {
@@ -10,6 +10,7 @@ import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.msg.TestTurnTokens;
import dev.ltms.bridged.peer.PeerUnreachableException;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
@@ -101,7 +102,7 @@ class SessionManagerTest {
assertDoesNotThrow(() -> sessions.asPresence().markPresent(null),
"the primary's null terminal must not blow up an unrelated tool call");
assertDoesNotThrow(() -> sessions.onDelivered(null));
assertDoesNotThrow(() -> sessions.onDelivered(null, TestTurnTokens.inert(null)));
assertDoesNotThrow(() -> sessions.onTurnComplete(null));
assertDoesNotThrow(() -> sessions.onTurnFailed(null));
@@ -121,7 +122,7 @@ class SessionManagerTest {
"MCP presence moves SPAWNING → READY");
assertTrue(sessions.asPresence().isPresent(terminal), "presence is also recorded");
sessions.onDelivered(terminal);
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
assertEquals(MemberSession.State.BUSY, sessions.get(session.paneId()).orElseThrow().state(),
"delivery moves READY → BUSY");
@@ -153,7 +154,7 @@ class SessionManagerTest {
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
sessions.onDelivered(terminal);
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
sessions.onTurnFailed(terminal);
@@ -181,7 +182,7 @@ class SessionManagerTest {
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
sessions.onDelivered(terminal);
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
sessions.onTurnFailed(terminal);
@@ -264,7 +265,7 @@ class SessionManagerTest {
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
sessions.onDelivered(terminal);
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
clock[0] = 100;
assertEquals(0, sessions.reapIdle(10), "BUSY session past TTL is never reaped");
@@ -281,7 +282,7 @@ class SessionManagerTest {
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
sessions.onDelivered(terminal);
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
sessions.onTurnComplete(terminal);
clock[0] = 21;
@@ -299,7 +300,7 @@ class SessionManagerTest {
MemberSession busy = sessions.acquire("ltms-local", "/busy", "/caller", "owner2");
sessions.asPresence().markPresent(ready.terminalId());
sessions.asPresence().markPresent(busy.terminalId());
sessions.onDelivered(busy.terminalId());
sessions.onDelivered(busy.terminalId(), TestTurnTokens.inert(busy.terminalId()));
clock[0] = 50;
assertEquals(1, sessions.reapIdle(30), "only READY past TTL is reaped");
@@ -317,9 +318,9 @@ class SessionManagerTest {
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
sessions.onDelivered(terminal);
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
sessions.onTurnComplete(terminal);
sessions.onDelivered(terminal);
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
sessions.onTurnComplete(terminal);
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
@@ -337,13 +338,13 @@ class SessionManagerTest {
String terminal = session.terminalId();
sessions.asPresence().markPresent(terminal);
sessions.onDelivered(terminal);
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
sessions.onTurnComplete(terminal);
assertEquals(MemberSession.State.DONE,
sessions.get(session.paneId()).orElseThrow().state(),
"first turn completes without release");
sessions.onDelivered(terminal);
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal));
sessions.onTurnComplete(terminal);
assertTrue(sessions.get(session.paneId()).isEmpty(), "session released after cap reached");
@@ -359,7 +360,7 @@ class SessionManagerTest {
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId());
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
assertTrue(sessions.onTurnCompleteWithPostAction(session.terminalId()));
MemberSession updated = sessions.get(session.paneId()).orElseThrow();
@@ -373,7 +374,7 @@ class SessionManagerTest {
SessionManager sessions = sessionManager(herdr, () -> 0L, 1, true);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId());
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
assertFalse(sessions.hasPostTurnAction(session.terminalId()),
"a session at its cap will be released, not reset for reuse");
@@ -388,7 +389,7 @@ class SessionManagerTest {
SessionManager sessions = sessionManager(herdr, () -> 0L, 0, false);
MemberSession session = sessions.acquire("ltms-local", null, "/caller", "term_primary");
sessions.asPresence().markPresent(session.terminalId());
sessions.onDelivered(session.terminalId());
sessions.onDelivered(session.terminalId(), TestTurnTokens.inert(session.terminalId()));
sessions.onTurnComplete(session.terminalId());
@@ -406,7 +407,7 @@ class SessionManagerTest {
MemberSession busy = sessions.acquire("ltms-local", "/busy", "/caller", "ownerB");
sessions.asPresence().markPresent(ready.terminalId());
sessions.asPresence().markPresent(busy.terminalId());
sessions.onDelivered(busy.terminalId());
sessions.onDelivered(busy.terminalId(), TestTurnTokens.inert(busy.terminalId()));
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(100));
@@ -6,14 +6,21 @@ import dev.ltms.bridged.guard.SubscriptionGuard;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.FakeHerdr;
import dev.ltms.bridged.herdr.WorkspaceControl;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.LoggerContext;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.read.ListAppender;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.msg.TestTurnTokens;
import dev.ltms.bridged.peer.MemberRole;
import org.junit.jupiter.api.Test;
import org.slf4j.LoggerFactory;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import static org.junit.jupiter.api.Assertions.*;
@@ -178,6 +185,77 @@ class WorktreeSessionManagerTest {
assertTrue(sessions.get(paneId).isEmpty(), "released session is no longer retrievable");
}
/**
* CB-576. A normal {@code COMPLETED} release whose worktree holds uncommitted work must NOT
* remove it — {@code --force} would destroy the worker's only copy. The bridge cannot see
* uncommitted files, so the worktree is preserved and the release logged at WARN naming the
* path, the session, and the cause an operator needs to find the work.
*/
@Test
void releasePreservesDirtyWorktreeAndLogsWarn() {
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
.withDirty(true);
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-576", null));
LoggerContext ctx = (LoggerContext) LoggerFactory.getILoggerFactory();
ch.qos.logback.classic.Logger sessionLog =
(ch.qos.logback.classic.Logger) LoggerFactory.getLogger(SessionManager.class);
ListAppender<ILoggingEvent> appender = new ListAppender<>();
appender.setContext(ctx);
appender.start();
sessionLog.addAppender(appender);
sessionLog.setLevel(Level.WARN);
try {
sessions.release(s.paneId());
assertTrue(herdr.called("pane.close"), "release still tears the worker pane down");
assertTrue(worktrees.removeCalls().isEmpty(),
"a dirty worktree is never removed — it holds the only copy of the work");
String warn = appender.list.stream()
.filter(e -> e.getLevel().equals(Level.WARN))
.map(ILoggingEvent::getFormattedMessage)
.filter(m -> m.contains("dirty worktree"))
.findFirst()
.orElse("no dirty-release WARN logged");
assertTrue(warn.contains(s.worktree()), "the WARN names the worktree path: " + warn);
assertTrue(warn.contains(s.terminalId()), "the WARN names the session: " + warn);
assertTrue(warn.contains("COMPLETED"), "the WARN names the release cause: " + warn);
} finally {
sessionLog.detachAppender(appender);
}
}
/**
* CB-576 review. A worktree that is already gone (operator cleanup, {@code git worktree prune},
* an earlier half-completed release) must not break teardown. {@code hasUncommitted} reports the
* missing path clean, so release still runs {@code notifyReleased} (the CB-516 fast-fail for a
* blocked {@code bridge_send} caller) and {@code launcher.stop} (so the pane is not orphaned),
* and falls through to the already-gone-tolerant {@code remove}.
*/
@Test
void releaseStillStopsPaneAndNotifiesWhenWorktreeIsGone() {
FakeHerdr herdr = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
MemberSession s = sessions.acquire("ltms-local", null, "/caller/proj", null,
new WorktreeRequest("cb-576g", null));
AtomicReference<String> releasedTerminal = new AtomicReference<>();
sessions.onRelease(releasedTerminal::set);
worktrees.markGone(s.worktree());
sessions.release(s.paneId());
assertEquals(s.terminalId(), releasedTerminal.get(),
"notifyReleased must still fire when the worktree is already gone (CB-516)");
assertTrue(herdr.called("pane.close"),
"the pane must still be stopped when the worktree is already gone");
assertEquals(1, worktrees.removeCalls().size(),
"release still calls the already-gone-tolerant remove");
}
@Test
void drainAllPreservesWorktreeOfIdleSession() {
FakeHerdr herdr = new FakeHerdr();
@@ -204,7 +282,7 @@ class WorktreeSessionManagerTest {
new WorktreeRequest("cb-544", null));
String terminal = s.terminalId();
sessions.asPresence().markPresent(terminal);
sessions.onDelivered(terminal); // BUSY, never completes → still BUSY when the timeout hits
sessions.onDelivered(terminal, TestTurnTokens.inert(terminal)); // BUSY, never completes → still BUSY when the timeout hits
sessions.drainAll(TimeUnit.MILLISECONDS.toNanos(100));
+4 -1
View File
@@ -16,6 +16,9 @@
AuditLogTest, which attaches its own ListAppender and asserts on emitted records.
-->
<!-- Keep the test logger behaviour aligned with the production cancellation filter. -->
<turboFilter class="dev.ltms.bridged.logging.McpCancelledNotificationFilter"/>
<appender name="STDOUT" class="ch.qos.logback.core.ConsoleAppender">
<encoder>
<pattern>%d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n</pattern>
@@ -34,4 +37,4 @@
<appender-ref ref="STDOUT"/>
</root>
</configuration>
</configuration>
+937
View File
@@ -0,0 +1,937 @@
# M4 - Fleet health, recovery, routing, and capacity
**Status:** Design accepted on 2026-08-15. CB-573 part 1 has shipped the classification model and
the `bridge_list` capacity view; the remaining M4 units are not yet shipped.
**Scope:** Fleet evidence, safe mechanical repair, lead routing, capacity reporting, and optional
human notification.
**Grounded in:** `health/FleetHealth`, `health/PaneBudget`, `inject/StatusPoller`,
`inject/StatusRefiner`, `inject/CompletionResolver`, `inject/Injector`, `session/SessionManager`,
`msg/MessageService`, `msg/ReplyInbox`, `msg/ReplyPushLoop`, `msg/LeadHeartbeatLoop`,
`mcp/PrimaryRegistry`, and `herdr/AgentControl`.
## 1. Problem and decision boundary
The operator asked the bridge to detect idle agents, exceptions, stopped work, and broken
communication. The bridge may read an agent pane from time to time. It must notify a person when
the fleet cannot move forward.
The four operator terms are not four equal health states. `IDLE` is a normal mode. An exception is
sometimes visible only as pane text. Stopped work may look the same as slow work. Broken
communication can occur on several links.
M4 uses this boundary:
- The bridge detects facts and joins evidence.
- The bridge repairs only mechanical failures with no judgement.
- The lead decides whether to stop, retry, replace, or reassign a member.
- A human is notified only when no healthy lead can act.
- n8n may route an outbound incident. It never classifies state or chooses recovery.
An inbound n8n decider would need bridge authority. No narrow machine-decider role exists. Giving a
workflow engine lead authority is unsafe, while adding a new role is a separate authorization
design. An outbound sink needs no bridge role.
The bridge must never replay a delivered task. That task may already have changed files, pushed a
branch, opened a pull request, or changed external state. A replay can run those side effects twice.
This rule must remain true even if later code stores delivered prompt text.
## 2. Evidence model
A health state is mainly a comparison between two views:
- **herdr view:** current agents and raw live status from one `AgentControl.list()` call.
- **bridge view:** session FSM, MCP presence, accepted turns, tasks, inbox state, and lead ownership.
A strong fault often appears as a disagreement between those views. For example, `BUSY` in the
session FSM and `DONE` in herdr means the bridge missed a turn boundary. Pane reads support this
model, but they are not the main monitor.
`SessionManager.rosterView` already joins session state and live status. `AgentControl.list()`
already gets the whole live fleet in one call. M4 makes that join persistent and adds timers,
accepted-turn state, and incident state.
### 2.1 Real traces behind the design
The first trace was an architect that stopped making progress:
```text
profile=opus role=architect state=busy liveStatus=done
```
The session moved from `DONE` to `BUSY` for turn 2. Eighteen minutes later, the session still said
`BUSY`, herdr still said `DONE`, the async task still said `PENDING`, and no completion fallback had
run. This is `TURN_BOUNDARY_LOST`, not a general slow-turn guess.
The second trace had two async sends to the same pane, one second apart. The pane was then stopped.
One ticket became failed. The other stayed `pending - worker unknown`. Current
`MessageService.abandon` resolves only `Rendezvous.currentWaiter(target)`, while async tasks live in
a separate ticket map. CB-568 is intended to fix that bug. M4 still keeps an independent
post-teardown invariant so a later regression becomes `DELEGATION_ORPHANED`.
### 2.2 Corrections made during design
The first state table missed `BUSY` in bridged plus `IDLE` or `DONE` in herdr. It would have found
the real trace only through a late, weak stall timer. The final model adds
`TURN_BOUNDARY_LOST` as a strong disagreement state.
The first notification design also required a webhook before `health.enabled` could turn on. That
removed useful local detection to avoid a narrower human-notification gap. The final design splits
detection from notification. Missing human escalation is shown as partial coverage instead of
disabling health.
## 3. Classification precedence
Evidence is applied in this order. A lower rule cannot hide a higher one.
1. **Control link:** failed fleet list plus failed ping becomes `CONTROL_LINK_DOWN`.
2. **Definitive target loss:** `_not_found` becomes `GONE` or `LEAD_UNREACHABLE` when the control
link is healthy.
3. **Startup and teardown invariants:** readiness expiry becomes `NEVER_READY`; surviving tasks
after teardown become `DELEGATION_ORPHANED`.
4. **Bridge/live disagreement:** `BUSY` plus stable raw `IDLE` or `DONE` becomes
`TURN_BOUNDARY_LOST`.
5. **Known screen evidence:** a tested fatal signature becomes `ERROR_ON_SCREEN`.
6. **Timed suspicion:** unchanged sparse pane probes may become `STALL_SUSPECTED`.
7. **Communication quality:** completion fallback becomes `MUTE`; an old inbox entry becomes
`REPLY_STRANDED`.
8. **Normal mode:** `STARTING`, `IDLE`, `WORKING`, `WORK_PENDING`, or `BLOCKED_AMBIGUOUS`.
The member flow in Figure 1 shows lifecycle states and the main fault exits. Fault states are
reported beside the session FSM; most are not new FSM values.
```mermaid
flowchart TD
Registered["Member registered"] --> Starting["STARTING"]
Starting -->|"MCP presence"| Idle["IDLE"]
Starting -->|"Readiness grace expires"| NeverReady["NEVER_READY"]
Idle -->|"Accepted delivery"| Working["WORKING"]
Working -->|"Trusted turn boundary"| Idle
Working -->|"Bridge BUSY and herdr IDLE or DONE"| Lost["TURN_BOUNDARY_LOST"]
Working -->|"Known fatal screen"| Error["ERROR_ON_SCREEN"]
Working -->|"Long age and unchanged sparse probes"| Stall["STALL_SUSPECTED"]
Working -->|"Target not found"| Gone["GONE"]
Idle -->|"Inbox or queued delivery exists"| Pending["WORK_PENDING"]
Pending -->|"Delivery or collection finishes"| Idle
Idle -->|"Raw BLOCKED with an open turn"| Blocked["BLOCKED_AMBIGUOUS"]
Lost -->|"Strict guarded repair"| Repaired["DONE with reconciled completion"]
Lost -->|"Repair refused"| LeadDecision["Lead decision required"]
```
*Figure 1. The member lifecycle and the main health exits. Pane-based states never authorise an
automatic retry of the task.*
## 4. State model
### 4.1 Normal and transitional member states
| State | Exact evidence | Meaning and certainty |
|---|---|---|
| `STARTING` | Session is `SPAWNING`; MCP presence is absent | Normal inside the startup grace. MCP contact is the readiness signal. |
| `IDLE` | Session is `READY` or `DONE`; live status is `IDLE` or `DONE`; no open turn or inbox item exists | Normal. Idle is not a fault. |
| `WORKING` | Session is `BUSY`; raw live status is `WORKING`; the accepted turn is open | Certain that herdr sees work. It does not prove useful progress. |
| `WORK_PENDING` | Queued delivery or inbox content exists while the target is injectable | Transitional. Existing injector or push logic should move it. |
| `BLOCKED_AMBIGUOUS` | An open turn exists and raw live status is `BLOCKED` | The bridge cannot tell whether this is permission, input, or a settled screen. |
Idle may drive configured resource cleanup. It never opens an incident and never pages a person.
### 4.2 Member fault and quality states
| State | Exact evidence | Certainty and action |
|---|---|---|
| `NEVER_READY` | `SPAWNING`, no MCP presence, and an accepted delivery waits through the existing readiness grace | Delivery never became possible. The exact cause is unknown. Fail the send, stop the process, and preserve a provisioned worktree. |
| `GONE` | Per-target herdr call returns `_not_found` while fleet list or ping works | Certain target loss. Fail all target work. Do not replay it. |
| `TURN_BOUNDARY_LOST` | Same session turn stays `BUSY`; same accepted task stays open; two raw snapshots show `IDLE` or `DONE` | Strong disagreement. Strict reconciliation may repair it. |
| `ERROR_ON_SCREEN` | Suspicious non-working state survives grace; `detection` matches a tested adapter-specific fatal signature | Certain only for the matched signature. A bare word such as `Exception` is not enough. |
| `STALL_SUSPECTED` | Open turn is older than the configured threshold; two normalised `recent_unwrapped` digests are unchanged; no boundary or reply occurs | Not certain. A long valid API call can look the same. Lead decides. |
| `MUTE` | Turn resolves through completion fallback instead of `bridge_reply` | Certain that no structured reply won. It does not prove an MCP failure. A single event is a metric, not an incident. |
| `REPLY_STRANDED` | Typed reply or health message remains after owning-lead push reaches its cap | Collection failed. This does not explain whether the lead is busy, dead, or ignoring the nudge. |
| `DELEGATION_ORPHANED` | Target is gone, failed, or released, but one or more tasks remain `PENDING` after reconciliation grace | Certain bridge invariant failure. This is not an inbox-drain fault. |
| `WORK_PRODUCT_AT_RISK` | Provisioned branch has commits after its recorded base; member is `DONE`, `FAILED`, or preserved after release; no turn or inbox item remains; long-idle threshold passed | A warning, not proof of loss. Work may already have an open pull request or a squash merge. |
`MUTE` opens an incident only after a small fixed rate threshold for one target or profile, or when
it appears with another fault.
`WORK_PRODUCT_AT_RISK` must not become `WORK_PRODUCT_UNCOLLECTED`. The bridge does not know pull
request or merge state. If committed work appears with `REPLY_STRANDED` or
`DELEGATION_ORPHANED`, the existing incident gains `committedWorkAtRisk: true`.
### 4.3 Control-link state
| State | Exact evidence | Certainty and action |
|---|---|---|
| `CONTROL_LINK_DOWN` | Two full-fleet `agent.list` calls fail across the grace, and herdr `ping` also fails | Certain for the bridged-to-herdr link. Retry calls, record the incident, and use human escalation if no lead can be reached. |
A failed fleet list alone is not a dead-member claim. A single `_not_found` with a healthy global
link is a target fault, not a control-link fault.
### 4.4 Lead states
| State | Exact evidence | Meaning and action |
|---|---|---|
| `LEAD_IDLE` | Expected lead is present with raw injectable status; no actionable state waits | Normal. Existing heartbeat may run under its own policy. |
| `LEAD_WORKING` | Expected lead is present with raw `WORKING`; stall threshold is not met | Reachable and busy. Never inject into the live turn. |
| `LEAD_STATUS_UNKNOWN` | Expected lead is present with raw `UNKNOWN` | Neither dead nor a healthy routing target. Retain evidence and retry. |
| `LEAD_UNREACHABLE` | Expected lead is absent from two successful live-agent snapshots while ping works, or targeted lookup returns `_not_found` with a healthy control link | Route to a healthy peer. If none exists, use human escalation. |
| `LEAD_UNRESPONSIVE` | Actionable state waits; lead stays injectable; bounded nudges exhaust; inbox remains uncollected | Route to a healthy peer or a person. |
| `LEAD_STALL_SUSPECTED` | Lead stays `WORKING` past threshold; two sparse pane probes show no progress | Not certain. Never kill or restart automatically. Route to peer or person. |
The monitor retains the lead name and terminal, last successful sighting, raw status and age,
consecutive list absences, targeted errors, pane-probe facts, pending incident age, and nudge
outcomes. Current heartbeat and push loops discard much of this history.
Expected lead identity comes from the same supplier used by `CallerResolver`. It is not liveness
evidence. `LeadTabScanner` keeps cached identity after a failed scan, so the health monitor compares
that identity with a fresh successful agent list. A dynamic identity also survives a two-successful-
snapshot retirement grace. This stops a dead lead from escaping health by disappearing from one map.
### 4.5 Evidence limits
M4 cannot tell these cases apart with current evidence:
- A valid long call and a hung call may have the same status and pane digest.
- `BLOCKED` does not explain which input is needed.
- An idle prompt after failure may look like an idle prompt after success.
- A missing structured reply does not prove a broken MCP connection.
- An undrained inbox does not explain why the lead did not collect it.
- Arbitrary pane text cannot safely classify arbitrary exceptions.
- A branch ahead of its base does not prove that work was not collected.
Logs are outputs, not classifier inputs. The monitor never parses its own logs.
## 5. Automatic action and lead action
### 5.1 Actions the bridge may take
The bridge may:
- retry transient herdr status, list, ping, and pane-read failures with bounded backoff;
- re-submit Enter after the existing paste/submit race;
- fail queued delivery after `NEVER_READY`;
- stop a never-ready process while preserving its provisioned worktree;
- fail all queued, accepted, and async tasks for a gone or released target;
- reconcile one lost boundary when every strict gate in Section 8 passes;
- hold typed messages, nudge the owning lead, and stop at the configured cap;
- use the existing bounded idle-lead heartbeat;
- deduplicate, route, update, and resolve incidents.
These actions do not choose new work and do not replay old work.
### 5.2 Decisions reserved for the lead
Only the lead may:
- stop or continue `BLOCKED_AMBIGUOUS`;
- stop, inspect, or wait on `ERROR_ON_SCREEN`;
- kill or continue `STALL_SUSPECTED`;
- spawn a replacement or reassign work;
- retry a delivered task;
- choose how to use partial work in a worktree;
- restart herdr or change network, model, credentials, backend, or configuration.
Reports include literal safe tool calls such as `bridge_status(sessionId="...")`,
`bridge_poll(ticket="...")`, `bridge_list()`, and optional `bridge_stop(paneId="...")`. A judgement
state never presents stop as the only action.
### 5.3 Release causes and worktree safety
| Release cause | Process action | Provisioned worktree |
|---|---|---|
| `SPAWN_ROLLBACK` before registration or delivery | Stop and clean up | Remove |
| `COMPLETED` for `READY` or `DONE` without pending work, idle TTL, or successful context-cap completion | Stop | Remove only if clean; preserve a dirty worktree (CB-576) |
| `NEVER_READY` | Stop | Preserve |
| `GONE` | Best-effort stop | Preserve |
| `TURN_FAILED` or lead abort while `BUSY` or `FAILED` | Stop | Preserve |
| `RELEASE_WITH_PENDING_TASKS` | Stop | Preserve |
| `SHUTDOWN` | Stop | Preserve |
Explicit stop is state-aware. `SPAWNING`, `BUSY`, `FAILED`, or any target with pending tasks uses a
preserving cause.
Before abnormal release removes the live session, M4 writes an atomic manifest under the worktree
root. It records session identity, owner, role, profile, repository, path, branch, base commit,
release cause, release time, state, and pending task ids. `bridge_list.preservedWorktrees` loads these
manifests after restart. Stop output and WARN logs also name the path and cause. M4 never
auto-deletes a preserved worktree.
## 6. Fleet health monitor
Add `FleetHealthMonitor`. Do not widen `StatusPoller` into a policy loop.
`StatusPoller` has a 250 ms delivery cadence and samples only injector targets with outstanding
work. Health needs all sessions, all leads, task state, inbox age, and global control evidence. One
loop cannot serve both cadences safely.
Build the monitor like `LeadHeartbeatLoop`:
- pure `decide(snapshot, priorState, now)` logic;
- a thin scheduler;
- an injected clock;
- edge-triggered state changes;
- no network work in the pure function;
- no sleeping in tests.
Each enabled fleet tick reads:
- one `AgentControl.list()` result for the whole fleet;
- one in-memory `SessionManager.roster()` snapshot;
- accepted turns and async task state;
- typed inbox depth, kind, and age;
- push and heartbeat outcomes;
- configured and discovered leads.
Existing failure paths publish structured evidence to the monitor. The monitor does not infer events
from log text.
### 6.1 Pane budget
Healthy idle members, recent working members, and quiet leads cause no pane reads.
A pane is eligible only for a stable lost boundary, sustained `BLOCKED` or `UNKNOWN`, work older
than the suspect threshold, or one final evidence read for a confirmed fault when the pane exists.
Compiled brakes apply even if config asks for more:
- per-target pane cooldown is at least 60 seconds;
- working age before the first progress probe is at least 300 seconds;
- at most two pane reads occur in one fleet tick;
- targets rotate fairly;
- only a normalised digest and optional clipped local excerpt are stored;
- no pane excerpt leaves bridged in a human webhook.
Use `detection` for tested screen signatures. Use normalised `recent_unwrapped` only for progress
comparison.
## 7. Typed inbox and routing
### 7.1 Semantic record
The typed inbox record carries:
```text
schemaVersion
kind: reply | health
msgId, target, subjectTerminal, recipientLead
severity, state, evidence
createdAtEpochMillis, firstSeenEpochMillis, lastSeenEpochMillis
recoveryTried, suggestedToolCalls, content
```
A health message never calls `Rendezvous.resolve`. It cannot look like the member's task result.
Both inbox adapters share field preservation, first-id-wins dedup, FIFO among decoded messages,
explicit ownership, ack, and release rules. The in-memory adapter stores typed records directly. It
does not copy AMQP migration logic.
### 7.2 AMQP migration
The reader uses AMQP `content_type`, never body sniffing:
```text
Legacy v0: text/plain
Typed family: application/vnd.ltms.bridged.inbox-message+json
```
A legacy reply may begin with `{`. It remains plain text because its media type is `text/plain`.
Legacy text becomes `kind=reply` with exact UTF-8 content and absent typed metadata.
Typed JSON has required integer `schemaVersion: 1`. Version 1 ignores unknown optional fields.
Missing required fields, invalid enums, malformed UTF-8 or JSON, and property/body identity mismatch
are invalid data.
An unknown schema version is not partly decoded. It remains unacknowledged on the original queue and
creates one operator-visible `unsupported_version` failure. A newer daemon may read it later.
Invalid known-format data is copied byte-for-byte to durable queue
`agent.<target>.inbox.quarantine`. A dedicated confirm-mode publisher confirms the persistent copy
before the original is acknowledged. A failed quarantine handoff leaves the original unacknowledged.
The raw body never enters logs.
Decode failure creates a redacted WARN, metric, `bridge_list` summary, and routed health incident.
One bad entry never escapes the consumer callback and never stops later valid messages.
Safe downgrade is not supported. The previous build ignores `content_type` and would show typed JSON
as ordinary reply text. If drained, it would acknowledge the message and lose typed meaning. Typed
queues must be drained or preserved before an old jar runs.
The existing contract suite uses RabbitMQ. Production uses LavinMQ. The migration and lead-key
ownership cases must run once against production LavinMQ before release, or the release must state
that LavinMQ was not checked.
### 7.3 Member routing
A member incident first goes to the exact lead that owns its accepted delegation.
`PrimaryRegistry` needs a no-fallback `delegatingLeadFor(memberTarget)` query. Health routing must not
use the old singular-primary fallback when several leads exist.
Publish the incident under the affected member target. Trigger the existing bounded push route. The
push waits until the owning lead is injectable, so it does not interrupt a live lead turn.
### 7.4 Peer lead routing
Figure 2 shows the route from incident to lead, peer, or person.
```mermaid
flowchart TD
Incident["Open incident"] --> Member{"Member incident?"}
Member -->|"yes"| Known{"Exact delegation owner known?"}
Known -->|"no"| Sink{"Human webhook enabled and healthy?"}
Known -->|"yes"| Owner{"Owner lead healthy?"}
Owner -->|"yes"| OwnerInbox["Publish to owner lead path"]
Owner -->|"no"| PeerSet["Build healthy peer candidate set"]
Member -->|"no, lead incident"| PeerSet
PeerSet --> Peer{"Healthy peer exists?"}
Peer -->|"yes"| Select["Choose fewest assigned incidents<br/>then stable name and terminal id"]
Select --> PeerInbox["Publish to peer lead inbox<br/>and status-gated push"]
Peer -->|"no"| Sink
Sink -->|"yes"| Webhook["Send classified outbound incident"]
Sink -->|"no"| Passive["Keep incident open<br/>show partial coverage on local surfaces"]
```
*Figure 2. Routing keeps delegation ownership separate from temporary peer fallback.*
Peer candidates exclude the incident subject, failed owner, absent leads, raw-unknown leads, and
leads with an open unhealthy state. A reachable `WORKING` peer may be selected; its push waits for an
injectable window.
Choose the candidate with the fewest assigned foreign incidents. Break ties by stable lead name,
then terminal id. Pin the recipient. Reassign only if that peer becomes unhealthy or retires. A
routing generation marks a reassignment, and old pending assignments become superseded.
`bridge_list` lead rows show health, health age, assigned foreign incident count, and a bounded list
of incident id, subject, state, severity, age, and routing generation. The top-level view also shows
owner, recipient, and routing reason.
A peer incident is published under the recipient lead's inbox key, not the failed subject's key. Its
status-gated nudge names the failed lead and gives the exact
`bridge_poll(target="<recipient-terminal>")` call.
### 7.5 Lead inbox ownership
Add `LeadInboxRegistry`, driven by the same expected-lead supplier as `CallerResolver`.
It calls `replyInbox.own(leadTerminal)` at startup for configured leads, after successful discovery,
after config adds a lead, and before publication. Ownership is not an authorization side effect.
A missing lead keeps its key owned. Release happens only after confirmed retirement, all incidents
are reassigned or resolved, typed health messages move or ack, and the queue is empty. Own a
replacement terminal before moving messages from the old key. Never release a non-empty in-memory
lead key, because in-memory release clears local data.
### 7.6 Single-lead deployment
One lead and no peer is a normal mode, not an edge case.
An idle, reachable lead may receive the existing bounded nudge. An unreachable or stalled sole lead
has no safe in-loop recovery. The bridge must not restart or replace it. A new lead would not have the
failed lead's plan or context, and an uncertain relaunch could create two orchestrators.
With no webhook, only `bridge_list`, `/healthz`, metrics, WARN logs, and the incident journal remain.
These are passive surfaces. They are not a human notification.
## 8. Lost-boundary reconciliation
This is the only M4 path that reconstructs a result. It must prefer a visible stall over a fabricated
reply.
### 8.1 Why normal completion rules are not enough
Current `CompletionResolver.resolve` has two fail-open rules. It resolves when the delivery baseline
is missing. It also resolves an empty completion when the pane read fails. Those choices are valid
after a trusted `WORKING -> IDLE` boundary because the bridge knows the turn ran. They are unsafe
when health only guesses that a boundary was lost.
M4 gives each accepted send an internal `TurnToken`. It ties target, exact waiter, session turn,
delivery baseline, and task outcome together.
### 8.2 Delivery baseline
Capture the baseline immediately after prompt send and before the delivery future completes. Store:
```text
TurnToken
exact waiter identity
capture time and pane source
normalised assistant block clipped to MAX_SCRAPE_CHARS
whether a supported assistant marker was recognised
capture result: PRESENT | READ_FAILED | UNRECOGNISED
```
A failed or missing baseline never authorises repair. A late baseline is not valid evidence. After a
daemon restart, the old waiter, task, token, and baseline are gone, so the old turn cannot be
repaired.
Automatic repair is enabled only for agent kinds with tested assistant-block fixtures. Current
extraction is Claude Code-specific and falls back to arbitrary raw text without `⏺`. That raw fallback
cannot authorise repair. OpenCode repair stays disabled until live pane fixtures exist.
### 8.3 Strict gates and resolver result
Figure 3 shows the repair gates. Any failed gate keeps the waiter unchanged.
The two raw snapshots must describe the same `TurnToken` and session turn. No `WORKING`,
`BLOCKED`, `UNKNOWN`, missing-agent, reply, failure, or new-delivery observation may occur between
them.
```mermaid
flowchart TD
Candidate["TURN_BOUNDARY_LOST candidate"] --> Stable{"Same TurnToken and BUSY turn<br/>across two raw IDLE or DONE snapshots?"}
Stable -->|"no"| Resnapshot["Take a fresh snapshot"]
Stable -->|"yes"| Waiter{"Exact captured waiter<br/>still open by identity?"}
Waiter -->|"no"| Stale["STALE_TURN or ALREADY_RESOLVED"]
Waiter -->|"yes"| Baseline{"Successful recognised<br/>delivery baseline exists?"}
Baseline -->|"no"| Refuse["Refuse repair<br/>leave ticket pending"]
Baseline -->|"yes"| Read{"Fresh pane read succeeds?"}
Read -->|"no"| Refuse
Read -->|"yes"| Output{"Recognised non-blank assistant block<br/>differs from clipped baseline?"}
Output -->|"no"| Refuse
Output -->|"yes"| Resolve["Shared CompletionResolver guard core<br/>resolves exact waiter"]
Resolve -->|"won race"| Repaired["RECONCILED_COMPLETION<br/>same turn becomes DONE"]
Resolve -->|"lost race"| Resnapshot
```
*Figure 3. Repair needs stronger evidence than a normal observed turn boundary.*
Refactor the current resolver into one guard core with two policies:
```text
resolveCaptured(target, inFlight, OBSERVED_BOUNDARY)
resolveCaptured(target, inFlight, LOST_BOUNDARY_REPAIR)
```
The health monitor calls only:
```text
CompletionResolver.reconcileLostBoundary(target, expectedTurnToken)
```
It returns `REPAIRED`, `ALREADY_RESOLVED`, `REFUSED_NO_CAPTURE`, `REFUSED_NO_BASELINE`,
`REFUSED_UNREADABLE`, `REFUSED_UNCHANGED`, `REFUSED_AMBIGUOUS_OUTPUT`, `STALE_TURN`, or
`RACE_LOST`.
Only `REPAIRED` and same-turn `ALREADY_RESOLVED` may move that turn from `BUSY` to `DONE`. A
per-target reconciliation gate stops a queued second send from being accepted between waiter
resolution and the FSM transition.
### 8.4 Lead-visible marker and refusal
A repaired result uses distinct `RECONCILED_COMPLETION` values in `Rendezvous`, `MessageService`,
task poll source, and metrics. The lead sees:
```text
[repaired completion - bridged detected a lost turn boundary. The member did not call
bridge_reply; pane-derived text follows and may be partial]
```
Clipped text also keeps the existing clipped-tail marker.
A refused repair leaves `TURN_BOUNDARY_LOST` open and the ticket pending. The report states that no
reply was reconstructed and no task was replayed. `UNCHANGED`, `UNREADABLE`, and
`AMBIGUOUS_OUTPUT` get at most one delayed retry for the same token. Missing capture or baseline gets
no retry. After two refused scrapes, automatic repair stops for that token.
### 8.5 Target-wide teardown invariant
CB-568 owns the multi-ticket cancellation mechanism. M4 routes every terminal cause through that one
idempotent operation and checks this independent invariant after teardown:
- no injector entry exists for the target;
- no accepted turn or completion record exists;
- no rendezvous waiter or ask exists;
- every async task is terminal or was already terminal;
- no thread waiting for the target send lock can later accept it;
- new sends fail immediately;
- each old task has one terminal outcome and one metric count.
A violation becomes `DELEGATION_ORPHANED`. The monitor may call the same idempotent target-wide
failure operation once. It never recreates the task.
## 9. Capacity and utilisation
Capacity is a view, not a health state.
`bridge_list` adds one block per profile:
```text
profile, maxLoad, live, free, reclaimable
```
For an unlimited profile, `maxLoad` and `free` are null. `free` is
`max(0, maxLoad - live)` for a capped profile.
The view must use the exact live-count function used by placement. A second calculation could show a
free slot that placement then refuses. Member rows add `idleForSeconds` only when state is `READY` or
`DONE`, no accepted turn exists, and the inbox is empty. `reclaimable` means only that the member
holds capacity without open bridge work.
The existing idle-lead nudge gains a bounded capacity summary. It lists per-profile live, cap, free,
and reclaimable counts, plus at most three long-idle members. Capacity does not make
`FleetState.hasPending()` true. A changed capacity fingerprint may re-arm one capped heartbeat
sequence. The fingerprint excludes changing idle durations, so a static idle fleet cannot reset the
cap forever. Reply-push stand-down remains first.
The bridge must never:
- spawn a member because a slot is free;
- generate a task or acceptance criteria;
- move queued work to another member or profile;
- treat a free slot or idle member as an incident;
- stop an idle member only to improve utilisation.
The bridge knows capacity facts but has no work list. Only the lead has the plan, task context,
side-effect history, and acceptance criteria.
Capacity calculation is in memory and adds no pane reads. Work-product checks run on a terminal
session edge, not every fleet tick.
This capacity design adds no automatic stop. The accepted `NEVER_READY` cleanup can still stop a
very slow startup after the existing grace, which is a known risk. Free capacity and long idle time
never trigger that path.
## 10. Human escalation and notification
### 10.1 Escalation rule
Notify a person only when no healthy lead can act:
- `CONTROL_LINK_DOWN` survives grace;
- a lead is unhealthy and no healthy peer can receive the incident;
- a member incident has no known owning lead;
- the only owning lead becomes unreachable, unresponsive, or stalled;
- incident publication or routing itself fails.
Do not page a person for a member fault while a healthy owning lead exists. An uncollected member
incident feeds lead-health evidence. If the lead then becomes unhealthy, peer or human routing starts.
### 10.2 Detection and notification switches
`health.enabled` controls detection and bridge-local reporting. It does not require a webhook.
`health.notifications.mode` is `disabled` or `webhook`. Disabled is valid and is the default.
Webhook mode requires a resolved environment variable. Turning notification off stops outbound
attempts but keeps incidents. Turning it back on resumes still-open human incidents.
Without a sink, `bridge_list.healthCoverage` states that human escalation is unavailable. `/healthz`
keeps its existing HTTP liveness result and adds a nested `fleetHealth.status=partial` component.
Metrics and one startup or reload WARN expose the same limit.
### 10.3 Incident and delivery deduplication
One open incident uses this key:
```text
(scope, subjectStableId, state, causeFingerprint)
```
The cause fingerprint includes stable error codes, dependency names, signature ids, or invariant
names. It excludes times, ages, retry counts, pane text, and changing digests. A later recurrence
after resolution gets a new generation and incident id.
Each outbound event uses:
```text
Idempotency-Key = hash(incidentId, eventType, eventRevision)
```
Event types are `open`, `severity_changed`, `reminder`, and `resolved`. Transport retries keep the
same key.
An atomic owner-only journal beside the active config stores open incidents, routing, delivered
revisions, retry state, and resolution state. It stores no pane or task content. Journal failure does
not stop detection, but notification coverage becomes degraded.
### 10.4 Retry, reminder, and resolve
Send the first event immediately. Retry network errors, timeouts, HTTP 408, HTTP 429, and HTTP 5xx
with full-jitter exponential backoff:
```text
base: 5 seconds
factor: 3
maximum delay: 15 minutes
one outstanding attempt per event
```
Respect `Retry-After` up to 15 minutes. Other HTTP 4xx responses are permanent for that event until
config changes or a person requests replay.
Transport retry is not an incident reminder. `humanRepeatSeconds` creates a new reminder revision
for an unresolved critical incident after the last successful human event. Disabled mode does not
build an unbounded reminder queue.
Send `resolved` only if at least one human event for that incident was delivered. If an incident
resolves before its first successful delivery, cancel the pending open event and record local
resolution.
### 10.5 Outbound payload boundary
An outbound payload may contain incident id and event type, severity, state, scope, stable bridge
ids, role or profile, times, duration, structured evidence type and counts, recovery attempted,
routing reason, coverage, and safe tool calls.
It must never contain:
- raw pane text, pane excerpts, or pane digests;
- task briefs, prompts, or member reply content;
- source files, diffs, or worktree file content;
- worktree paths;
- environment values, tokens, credentials, headers, or webhook URL;
- raw exception messages or stack traces;
- arbitrary model output.
The sink response body is ignored. A webhook cannot direct recovery. n8n remains outbound-only.
### 10.6 Metrics
M4 adds bounded-label series:
```text
bridged_health_incidents{scope,state,severity}
bridged_health_incidents_total{event}
bridged_health_notifications_total{event,outcome}
bridged_health_notification_queue_depth
bridged_health_notification_last_success_seconds
bridged_health_notification_capability{mode,status}
bridged_lead_health{lead,state}
bridged_lead_assigned_incidents{lead}
```
Metric labels never include terminal ids, incident ids, URLs, or error text.
## 11. Configuration
The optional `health:` block is absent or disabled by default. The dormant monitor scheduler does no
herdr or pane work while disabled. Every listed key is hot because the monitor reads `ConfigRef` on
each tick or notification.
| Key | Class | Default and hard bound | Purpose |
|---|---|---|---|
| `health.enabled` | Hot | `false` | Enable detection and bridge-local reporting. |
| `health.snapshotIntervalSeconds` | Hot | default 30, minimum 15 | Whole-fleet comparison cadence. |
| `health.workingSuspectAfterSeconds` | Hot | default 600, minimum 300 | Age before working-pane probes. |
| `health.paneProbeIntervalSeconds` | Hot | default 60, minimum 60 | Per-target pane cooldown. |
| `health.leadUnresponsiveAfterSeconds` | Hot | default 300, minimum 120 | Delay after exhausted actionable nudges before lead fault. |
| `health.humanRepeatSeconds` | Hot | default 3600, minimum 900 | Minimum repeat period for one open human incident. |
| `health.capacityLongIdleAfterSeconds` | Hot | default 900, minimum 300 | Long-idle threshold for capacity summaries. |
| `health.includePaneExcerpt` | Hot | `false` | Allow a clipped excerpt in local lead reports only. Human payloads still exclude it. |
| `health.notifications.mode` | Hot | `disabled` | Select `disabled` or `webhook`. |
| `health.notifications.webhookUrlEnv` | Hot | required in webhook mode | Name of the environment variable that holds the sink URL. |
| `health.notifications.requestTimeoutMs` | Hot | default 10000, range 1000-30000 | Whole webhook request limit. |
Two consecutive snapshots are compiled floors for lost boundary, lead disappearance, and control
link failure. The two-pane-reads-per-tick limit is also compiled and cannot be weakened by config.
## 12. Delivery units and acceptance
### Unit 1 - Evidence model and fleet snapshot
Scope: health state model, fleet join, clocks, evidence retention, and pane budget.
Acceptance criteria:
1. One `agent.list` call covers one enabled fleet tick.
2. Pure decision tests cover every state and every evidence limit in Section 4.
3. `BUSY` plus stable raw `DONE` opens `TURN_BOUNDARY_LOST` after two snapshots.
4. Healthy fleet snapshots perform zero pane reads.
5. Pane cooldown, two-read fleet budget, and fair rotation cannot be disabled by config.
6. Logs are outputs only; no log parsing exists.
7. Fleet snapshots expose the same profile live-count calculation that placement uses.
8. Capacity rows report cap, live, free, and reclaimable values without opening incidents.
### Unit 2 - Lost boundary and task reconciliation
Scope: accepted-turn identity, guarded repair, target-wide teardown, release causes, and preserved
worktree discovery.
Acceptance criteria:
1. Every accepted send receives a stable `TurnToken` tied to target, exact waiter, and delivery
baseline.
**Corrected during implementation (2026-08-15).** This criterion first also required the session
turn number and the task outcome. That is not implementable at this layer, and the implementer
refused it three times rather than fabricate a value — correctly. The reason is an ordering fact
that is invisible from any single class: `MessageService` owns acceptance and holds the waiter and
the async `Task`, but it learns nothing about delivery, because the delivery event goes to
`CompletionResolver` through `TurnListener.onDelivered`. And `CompletionResolver.onDelivered` runs
*before* `SessionManager.onDelivered`, so the session turn number does not exist yet at the only
point where the token could capture it.
Two ways out were rejected. A shared registry keyed by target reintroduces exactly the "whichever
send happens to be waiting" ambiguity the token exists to remove — the same weak claim
`Rendezvous.currentWaiter` warns about. Injecting a turn counter into `MessageService` adds a
required cross-layer dependency to populate a field that nothing in this slice reads, which is
speculative coupling across a boundary already shown to be fragile.
So the token identifies the **accepted send**, and `SessionManager` keeps verifying its own
delivery separately. Repair (criterion 2) does need the session turn; binding it means resolving
that acceptance-versus-delivery ordering first, and that work belongs to the repair unit, not
here. The token record carries a comment saying the field is deliberately absent.
2. Repair requires the same `BUSY` token, two raw `IDLE` or `DONE` snapshots, no conflicting
observation, exact open waiter, successful baseline, and new recognised assistant output.
3. Missing, failed, late, or post-restart baseline never authorises repair.
4. Repair is enabled only for agent kinds with tested assistant-block extraction. Raw-text fallback
without a recognised marker refuses repair.
5. Normal completion and repair use one resolver guard core. Waiter, scrape, clipping, unchanged, and
exact-turn guards are not duplicated.
6. `reconcileLostBoundary` returns every typed result named in Section 8.3.
7. Only `REPAIRED` and same-turn `ALREADY_RESOLVED` may move the same turn to `DONE`.
8. A per-target reconciliation gate blocks a queued second send during repair and FSM update.
9. Repaired completion has distinct rendezvous kind, message outcome, poll source, lead marker, and
metric. Clipping keeps its extra marker.
10. Unchanged, unreadable, or ambiguous evidence gets at most one delayed retry. Missing capture or
baseline gets none.
11. Refusal leaves the ticket pending and tells the lead that no result was rebuilt or replayed.
12. Release, gone, never-ready, and abnormal stop use CB-568's one idempotent target-wide failure
operation.
13. The post-teardown invariant in Section 8.5 is tested independently of CB-568 internals.
14. A violated teardown invariant creates `DELEGATION_ORPHANED` and retries only the idempotent
failure operation.
15. `SPAWN_ROLLBACK` and normal `COMPLETED` remove worktrees. Abnormal and shutdown causes preserve
them.
16. Explicit stop is state-aware. Any pending task or non-terminal state preserves the worktree.
17. Atomic preserved-worktree manifests reload after restart and appear in lead-only
`bridge_list.preservedWorktrees`.
18. Manifest failure preserves the worktree and opens an operator-visible health failure.
19. Provision records the base commit. Terminal, long-idle worktrees report
`WORK_PRODUCT_AT_RISK` only under the evidence in Section 4.2 and never auto-delete work.
20. No path replays a delivered task, rebuilds its brief, or retargets it, even when prompt text is
available.
21. Tests cover both real traces, all repair refusals, clipping, explicit-reply and next-turn races,
restart without capture, concurrent send and release, and preserved discovery after restart.
### Unit 3 - Typed inbox and member routing
Scope: semantic record, AMQP migration, both adapters, member routing, polling, and member health in
`bridge_list`.
Acceptance criteria:
1. AMQP selects legacy or typed decoding only from `content_type`; it never sniffs the body.
2. Persistent `text/plain` from the old build becomes `kind=reply` with exact UTF-8 content,
including content beginning with `{`.
3. New entries use the vendor media type, `schemaVersion: 1`, UTF-8, persistent delivery, and AMQP
message ids.
4. Version 1 ignores unknown optional fields but rejects missing fields and identity mismatch.
5. Unknown versions are not decoded or acked. They remain on the original queue and create one
deduplicated failure.
6. Invalid known data never escapes the callback, appears as a reply, or blocks later valid messages.
7. Invalid data reaches durable per-target quarantine before original ack. Failed handoff leaves the
original unacked.
8. Decode failures create redacted WARN, metric, `bridge_list` summary, and routed incident without
raw content.
9. Both adapters pass one semantic contract for fields, FIFO, dedup, ownership, ack, and release.
10. Lead keys require explicit ownership. Publication never claims a queue.
11. Unit codec tests cover legacy `{`, Unicode, malformed UTF-8, typed round trip, additive fields,
malformed JSON, missing fields, identity mismatch, media type, version, and dedup.
12. A live broker contract writes old wire data and reads it with the new adapter after reconnect.
13. Live contract tests cover mixed entries, quarantine confirm-before-ack, unsupported redelivery,
later progress past poison, property persistence, lead ownership, and ack removal.
14. Safe downgrade is documented as unsupported.
15. RabbitMQ contract tests pass with `mvn test -Pcontract`. The same cases run once on production
LavinMQ, or the release states that LavinMQ was not checked.
16. Member incidents route to the exact delegating lead and never resolve a task rendezvous.
17. `bridge_list` shows compact member health and capacity without pane content. Member
`idleForSeconds` is present only when no accepted turn or inbox item exists.
### Unit 4 - Lead health and peer routing
Scope: lead evidence, exact ownership, peer selection, explicit-recipient push, and lead inbox
lifecycle.
Acceptance criteria:
1. Lead identity uses the `CallerResolver` supplier. Liveness uses successful current agent data.
2. Two successful-list absences with healthy ping become `LEAD_UNREACHABLE`; global link failure does
not.
3. Raw `WORKING`, raw `UNKNOWN`, first-seen time, failures, last success, and error class persist
across ticks.
4. Heartbeat and push publish status and nudge outcomes before safe no-injection decisions.
5. Dynamic lead identity survives a two-successful-snapshot retirement grace.
6. Member incidents first use exact delegation ownership with no singular-primary fallback.
7. Peer selection follows the exclusions, load rule, and stable tie break in Section 7.4.
8. A selected working peer is not interrupted. Its push waits for an injectable window.
9. Recipient assignment stays pinned. Reassignment increments generation and supersedes old pending
assignment.
10. `bridge_list` shows bounded foreign assignments, recipient, reason, and generation without pane
content.
11. `LeadInboxRegistry` owns configured and discovered lead keys before publication.
12. Missing leads keep ownership. Retirement needs an empty queue and handled incidents.
13. Replacement owns the new key before messages move. Non-empty in-memory keys are not released.
14. Tests cover dead versus busy, unknown, global failure, stale scan cache, disappearance, one peer,
several peers, reassignment, and no peer.
15. Adapter tests cover lead ownership, restart re-ownership, retirement, and terminal replacement.
LavinMQ is checked or named as unchecked.
16. A sole unreachable or stalled lead is never restarted or replaced. Without a sink, only passive
evidence remains and every coverage surface says so.
### Unit 5 - Human sink, hot config, metrics, and operator coverage
Scope: generic webhook, config split, incident journal, retry, resolve, metrics, example config, and
operator documentation.
Acceptance criteria:
1. `health.enabled` works without a human sink.
2. Notification mode is hot, defaults to disabled, and supports disabled or webhook.
3. Webhook mode requires a resolved environment value. Bad notification config does not disable an
already valid detector.
4. Mode changes keep open incidents. Re-enable resumes eligible incidents.
5. `bridge_list`, `/healthz`, metrics, and one WARN show partial coverage without a sink. HTTP
liveness behavior stays unchanged.
6. One-lead, no-sink coverage states that lead failure has no active notification or recovery.
7. Incident and outbound dedupe use the stable keys in Section 10.3.
8. The owner-only local journal survives restart and contains no pane or task content.
9. Journal failure keeps detection running but marks notification coverage degraded.
10. Retry tests cover network failure, timeout, 408, 429, `Retry-After`, 5xx, permanent 4xx, jitter,
delay cap, config re-arm, and one outstanding attempt.
11. Reminders and transport retries remain separate. Disabled mode does not build an unbounded queue.
12. Resolve sends only after an earlier human event succeeded. Resolve-before-delivery cancels stale
open delivery.
13. Metrics use bounded labels and exclude ids, URLs, and error text.
14. Payload tests reject every content type forbidden in Section 10.5.
15. Webhook response bodies are ignored and cannot direct recovery.
16. Tests cover disabled mode, one lead without sink, open/update/reminder/resolve, restart, dedup,
reassignment, disable/re-enable, and sink failure while local health continues.
17. `bridged.example.yaml` documents all hot keys and compiled floors.
18. The operator Features wiki is updated separately. The portable `CLAUDE.md` block is checked and
changed only if shipped tool or inbox semantics make it untrue.
19. `mvn clean install` passes.
## 13. Not checked and release gates
These limits are part of the design, not optional follow-up notes.
- **OpenCode pane status and assistant markers were not checked.** OpenCode lost-boundary repair is
disabled until live fixtures exist.
- **Permission-prompt status was not checked** for Claude Code or OpenCode. `BLOCKED` remains
ambiguous and has no automatic action.
- **`recent_unwrapped` stability was not checked** across all supported agent kinds. If normalisation
is not stable, `STALL_SUSPECTED` must say its evidence is weaker.
- **The real `BUSY + DONE` trace was not replayed against live herdr.** The design uses the observed
production trace and current poller behavior.
- **CB-568 was not present when Unit 2 was designed.** Unit 2 must inspect the landed API and keep its
independent teardown invariant.
- **Production LavinMQ was not checked.** Existing durable-inbox contracts use RabbitMQ. Migration,
quarantine, redelivery, lead ownership, and reassignment must run on LavinMQ before release or be
recorded as unchecked.
- **Live multi-lead routing was not checked.** Peer choice and reassignment are design rules backed by
fake-clock and adapter tests until a live exercise runs.
- **A live sole-lead failure with a webhook was not checked.** The no-peer path is a design result,
not a tested recovery.
- **No n8n, Slack, PagerDuty, or other receiver was checked.** The webhook remains generic and
outbound-only.
- **Deployment supervisor behavior for nested `/healthz` fields was not checked.** HTTP liveness
status stays unchanged to reduce this risk.
- **Incident-journal crash behavior was not checked** because the journal does not exist yet. Unit 5
must test atomic replacement and restart recovery.
- **Worktree merge state cannot be checked reliably** without forge or explicit collection evidence.
`WORK_PRODUCT_AT_RISK` stays a warning.
## 14. Locked exclusions
M4 does not expose `agent.read` as a bridge tool. It does not add a workflow engine, inbound n8n
authority, automatic task assignment, task replay, automatic lead replacement, or automatic member
spawn for free capacity.
The bridge remains a message bus with evidence and bounded mechanical repair. The lead remains the
place where judgement and work planning happen.