Compare commits

..

8 Commits

Author SHA1 Message Date
Dai Ha 7b98cca967 LeadMailbox: durable leader-to-leader mailbox over a shared coordination vhost
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Failing after 1m35s
Adds the broker-side mechanism for lead-to-lead messages across daemons/hosts
(unit 1 of 2): a LeadMessage envelope carrying from/to coord-ids, an
AMQP-backed LeadMailbox modeled closely on AmqpReplyInbox (consume-and-hold,
deferred manual ack, confirm-mode publish, recovery handling), and a new
optional coordinator: config block (separate vhost from broker:, leader
traffic only). Config parsing + accessors only — FleetMcp/Fleetd/Injector/
MessageService and the send path are untouched; wiring is a separate ticket.
2026-08-24 17:23:42 +02:00
Dai Ha 7655f1b51a CB-634: pin the IDE overlay to the module dir + best-effort auto-open
The overlay pinned project_path to the worktree root. For a repo whose Maven
module is a subdir (this repo's pom is in `bridged/`, not at the root), opening
the root imports no module and every ide_* call resolves nothing. Pin and open
the module dir instead.

Two new opt-in per-Profile keys, both read only when ideMcpUrl is set:
- ideProjectDir: repo-relative module dir the IDE opens and the overlay pins;
  blank keeps the old worktree-root behaviour.
- ideOpenCommand: host command that opens that dir in the IDE at spawn, with
  {dir} substituted and run through /bin/sh -c so env (e.g. DISPLAY) can be set
  inline. Best-effort and non-fatal — a failure never fails the spawn. Blank
  keeps the manual-open behaviour. No close half yet (deferred).

Shared helpers PeerLauncher.ideProjectPath / openInIde back both launchers.
The two Profile fields ride a back-compat constructor, so every existing call
site and YAML compiles and behaves unchanged.

Tests: overlay content pins the module dir when ideProjectDir is set;
ideProjectPath resolution; openInIde no-op on a blank command. 918 tests green.
2026-08-24 06:47:37 +02:00
Dai Ha a5efb7c676 CB-634: write the overlay exclude to the common git dir, not the per-worktree gitdir
git reads info/exclude from the common dir for a linked worktree (only
info/sparse-checkout is per-worktree), so the entry written into
<common>/worktrees/<name>/info/exclude was never honoured and CLAUDE.local.md
showed as untracked -- at risk of being swept into a worker's PR. Derive the
common dir (<common>/worktrees/<name> -> <common>) and write there. Found by
dogfooding a real spawn on fleet01; the test now uses the real worktree layout
and asserts the entry lands in the common dir, not the per-worktree gitdir.
2026-08-23 20:06:36 +02:00
Dai Ha 5a8cf4cb4d Merge CB-634 overlay redesign: deliver IDE guidance as on-disk CLAUDE.local.md / opencode instructions overlay 2026-08-23 16:45:50 +02:00
Dai Ha a3fc7e4df8 CB-634: deliver IDE guidance as an on-disk overlay, not the system-prompt charter
Move the IDE guidance text to PeerLauncher.ideOverlayText (shared by both
launchers). ClaudeCodeLauncher drops it from the reply-charter file and writes
CLAUDE.local.md into a provisioned worktree instead, gated on a .git FILE
(safety: never writes into the primary's real .git-DIRECTORY checkout) and
registers it in info/exclude. OpenCodeLauncher mounts the intellij server and
adds the rules file to the instructions array.
2026-08-23 16:37:08 +02:00
Dai Ha d811b30df3 CB-634 (draft): mount the IDE Index MCP into a member, opt-in per profile
Adds `ideMcpUrl` to FleetConfig.Profile (default off). When set, the
Claude Code launcher mounts the IDE Index MCP as a second inline
--mcp-config server named `intellij`, and appends an IDE charter that
pins every ide_* call to the member's own worktree (spec.cwd()). The
charter order is role -> ide -> reply, one --append-system-prompt-file,
reply last (CB-618). The mount gate now fires on ideMcpUrl alone, not
only mcpUrl. ConfigRef treats an ideMcpUrl change as deferred, like the
other launch flags.

Never touches .mcp.json or CLAUDE.md — the mount and the rule arrive as
launch flags, so a project's own config is untouched.

Not yet done (see fleetd #162): the bridged-owned IDE lifecycle
(open on provision, close before worktree removal), and the opencode
adapter (separate ticket). fleetd.example.yaml documents ideMcpUrl and
fixes the stale parityOverlay default.

911 tests green.
2026-08-23 16:00:36 +02:00
Dai Ha 3f4ac2b24e CB-635: --check reports whether broker.uriEnv resolves in a login shell
CI / contract (push) Successful in 46s
CI / build (push) Failing after 1m30s
An empty uriEnv no longer stops the daemon (#152), so the failure is quiet: bridged
starts, falls back to the in-memory reply inbox, and held reports stop surviving a
restart. --check is the only thing that says so before the fact. The var name is read
out of bridged.yaml so a renamed key cannot make the check lie.
2026-08-23 14:03:08 +02:00
Dai Ha 4644359128 Merge CB-635: broker.uriEnv keeps the AMQP password out of the config, and an unreachable broker no longer stops the daemon (#151, #152)
CI / contract (push) Successful in 1m11s
CI / build (push) Failing after 1m42s
2026-08-23 13:58:21 +02:00
16 changed files with 1325 additions and 31 deletions
+37 -2
View File
@@ -126,6 +126,20 @@ herdrSocket: ~/.config/herdr/herdr.sock
# Use `pane` for the legacy behaviour (split the focused tab).
# mcpUrl → bridged mounts the bridge MCP (--mcp-config, inline) + reply charter
# (--append-system-prompt) as launch flags; nothing is written to the profile.
# ideMcpUrl → opt-in (CB-634), default off. When set, bridged mounts the IDE Index MCP as a
# second inline server named `intellij`, and adds an IDE charter that pins every
# ide_* call to the member's own worktree. A URL, not a boolean — host and port
# are host-specific. Set it only on a host where the IDE actually runs.
# ideProjectDir → repo-relative module dir the IDE opens and the overlay pins (CB-634). Only read
# when ideMcpUrl is set. This repo's Maven pom lives in `bridged/`, not at the
# worktree root, so opening the root imports no module and ide_* resolves nothing;
# set this to `bridged`. Omit for a repo whose project is the worktree root.
# ideOpenCommand → host command that opens ideProjectDir in the IDE at spawn (CB-634 auto-open).
# Only read when ideMcpUrl is set. `{dir}` is replaced with the absolute module
# dir and the command runs through `/bin/sh -c`, so set env inline if needed —
# e.g. `env DISPLAY=:10.0 idea {dir}`. Best-effort: a failure is logged, never
# fails the spawn. Omit to open the member's module by hand. There is no close
# half yet — an opened module stays open until the operator closes it.
# tokenEnv → host env var holding the worker's auth token (value never stored in config);
# omit for a backend that needs no token (e.g. a local ollama).
# cwd → pin this profile's working directory (CB-112). Omit to inherit the primary's
@@ -135,7 +149,8 @@ herdrSocket: ~/.config/herdr/herdr.sock
# skills/MCP/hooks. Omit to leave the worker on the host default.
# parityOverlay → repo-relative paths copied primary→worktree so a worker in a provisioned
# worktree sees the same local config (CB-301-ext). Omit for the default set:
# [.claude/settings.local.json, .env, .envrc].
# [.env, .envrc]. (.claude/settings.local.json is NOT in the default — it
# pre-approves IDE/tool grants a member must not hold ambiently; CB-525/CB-634.)
#
# Do NOT add .mcp.json (CB-525). A worker's tools are whatever its launcher
# mounts — the bridge, and nothing else. Replicating the primary's MCP config
@@ -237,7 +252,10 @@ profiles:
# credentialId: shared-openai # opt-in: quarantine together with every other profile sharing this id (CB-578)
# configDir: /Users/me/.ccs/instances/gx10 # CLAUDE_CONFIG_DIR — inherit that profile's skills/MCP
# cwd: /Users/me/src/myrepo # pin the working dir; omit to inherit the primary's
# parityOverlay: [".claude/settings.local.json", ".env", ".envrc"] # never add .mcp.json — see above
# parityOverlay: [".env", ".envrc"] # the default; never add .mcp.json or .claude/settings.local.json — see above
# ideMcpUrl: http://127.0.0.1:29170/index-mcp/streamable-http # opt-in (CB-634): IDE code intelligence, pinned to the worktree
# ideProjectDir: bridged # CB-634: module dir the IDE opens + the overlay pins (this repo's pom is in bridged/)
# ideOpenCommand: env DISPLAY=:10.0 idea {dir} # CB-634 auto-open: opens {dir} in the IDE at spawn; omit to open by hand
gx11: # a second backend, so `placement: weighted` has a choice
baseUrl: http://gx01.gw:8000 # self-hosted; ccs handles the model + token
placement: tab
@@ -636,6 +654,23 @@ guard:
# uriEnv: LAVINMQ_URI
# prefetch: 32
# Shared cross-host LEADER coordination broker. OMIT this block to leave lead-to-lead messaging
# off entirely (config-only in this ticket — nothing here wires it into a live LeadMailbox yet).
# This is a SEPARATE AMQP vhost from `broker:` above: member/worker inboxes always stay on the
# per-fleet `broker:` vhost, and this vhost carries only leader-to-leader traffic, so two fleets
# whose members must never see each other can still share one coordination vhost for their leads.
# uriEnv → name of a host env var holding the coordination AMQP URI, same convention as
# broker.uriEnv (keeps the credential out of fleetd.yaml). Wins over `uri` when set.
# selfId → this daemon's own lead coord-id — the name its mailbox is owned under
# (lead.<selfId>.inbox), e.g. "mac-opus" or "fleet01-lead". Must be globally unique
# across every daemon sharing this vhost.
# prefetch → consumer basicQos, capping how many unacked messages the mailbox holds in-heap.
# Default 32 when omitted.
# coordinator:
# uriEnv: LEAD_COORD_URI
# selfId: mac-opus
# prefetch: 32
# Active push-to-primary (CB-307 Stage 3). When a worker reply lands with no open fleet_send,
# the ReplyPushLoop injects a *drain nudge* (never the payload) into the primary's own herdr
# pane — status-gated (only when injectable, never mid-turn) and bounded. Ack = drain: the loop
@@ -276,6 +276,9 @@ public final class ConfigRef implements Supplier<FleetConfig> {
&& Objects.equals(a.workspace(), b.workspace())
&& Objects.equals(a.tabLabel(), b.tabLabel())
&& Objects.equals(a.mcpUrl(), b.mcpUrl())
// CB-634: the IDE MCP mount is a launch flag, fixed at spawn like mcpUrl — a
// reload changes it only for members spawned after, so a changed value is deferred.
&& Objects.equals(a.ideMcpUrl(), b.ideMcpUrl())
&& Objects.equals(a.cwd(), b.cwd())
&& Objects.equals(a.parityOverlay(), b.parityOverlay())
&& Objects.equals(a.gitTokenEnv(), b.gitTokenEnv())
@@ -6,6 +6,7 @@ import com.fasterxml.jackson.core.JsonToken;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
import dev.ltms.fleet.msg.AmqpReplyInbox;
import dev.ltms.fleet.msg.LeadMailbox;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.placement.PlacementPolicies;
import org.slf4j.Logger;
@@ -69,6 +70,11 @@ import java.util.Set;
* @param memberCredentials deny-by-default policy (CB-596) for which of the operator's own host
* credentials a spawned member's pane inherits. {@code null} (the block
* omitted) blocks nothing — see {@link MemberCredentials}.
* @param coordinator shared cross-host leader coordination broker: a SEPARATE AMQP vhost from
* {@link #broker} used only for lead-to-lead traffic (member/worker inboxes
* stay on {@code broker}'s vhost). {@code null} → no lead mailbox is opened.
* Config parsing + accessors only — nothing here wires it into a live
* {@code LeadMailbox}; that is a separate ticket. See {@link Coordinator}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record FleetConfig(
@@ -89,7 +95,20 @@ public record FleetConfig(
Auth auth,
ConfigReload configReload,
Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials) {
MemberCredentials memberCredentials,
Coordinator coordinator) {
/** Back-compat form before the {@code coordinator:} block was added. */
public FleetConfig(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, Health health, String placement, Auth auth,
ConfigReload configReload, Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, null);
}
/** Back-compat form before the CB-596 {@code memberCredentials:} block was added. */
public FleetConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
@@ -289,7 +308,10 @@ public record FleetConfig(
Integer maxLoad,
Boolean subscription,
String exhaustedPattern,
String credentialId) {
String credentialId,
String ideMcpUrl,
String ideProjectDir,
String ideOpenCommand) {
/** Peer kind spawned by {@link dev.ltms.fleet.member.ClaudeCodeLauncher} (the default). */
public static final String KIND_CLAUDE_CODE = "claude-code";
@@ -348,6 +370,17 @@ public record FleetConfig(
// "quarantines alone" fallback actually lives, so today's behaviour needs no defaulting
// here at all.
credentialId = (credentialId == null || credentialId.isBlank()) ? null : credentialId;
// CB-634: opt-in per profile, default off. A URL, not a boolean — host and port are
// host-specific, mirroring mcpUrl. When set, the member gets the IDE Index MCP mounted
// (pinned to its own worktree via the charter). Blank ⇒ off.
ideMcpUrl = (ideMcpUrl == null || ideMcpUrl.isBlank()) ? null : ideMcpUrl;
// CB-634 auto-open: both are only read when hasIdeMcp(). ideProjectDir is the repo-relative
// module dir IntelliJ must open (this repo's pom lives in `bridged/`, not at the root), and
// it is also the project_path the overlay pins. Blank ⇒ the worktree root (unchanged before
// auto-open). ideOpenCommand is the host command that opens that dir in the IDE, with {dir}
// substituted; blank ⇒ no auto-open (the operator opens the module by hand).
ideProjectDir = (ideProjectDir == null || ideProjectDir.isBlank()) ? null : ideProjectDir;
ideOpenCommand = (ideOpenCommand == null || ideOpenCommand.isBlank()) ? null : ideOpenCommand;
}
/**
@@ -392,7 +425,25 @@ public record FleetConfig(
public Profile withProfile(String p) {
return new Profile(p, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, subscription,
exhaustedPattern, credentialId);
exhaustedPattern, credentialId, ideMcpUrl, ideProjectDir, ideOpenCommand);
}
/**
* Backward-compatible constructor without the CB-634 auto-open fields
* ({@code ideProjectDir}/{@code ideOpenCommand}) — a profile that opts into IDE MCP still
* pins the worktree root and does not auto-open. Keeps pre-auto-open call sites (and any YAML
* that omits the keys) compiling and behaving identically. This is the shape the canonical
* constructor had before the two fields were added.
*/
public Profile(String profile, String baseUrl, String model,
String configDir, String tokenEnv, List<String> argv,
String placement, String workspace, String tabLabel, String mcpUrl,
String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv,
String kind, Map<String, String> env, Float weight, Integer maxLoad,
Boolean subscription, String exhaustedPattern, String credentialId, String ideMcpUrl) {
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad,
subscription, exhaustedPattern, credentialId, ideMcpUrl, null, null);
}
/** True when this profile is served by the Claude Code adapter (the default kind). */
@@ -441,7 +492,7 @@ public record FleetConfig(
Boolean subscription, String exhaustedPattern) {
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad,
subscription, exhaustedPattern, null);
subscription, exhaustedPattern, null, null);
}
/** True when this profile's CB-578 stage A backend-exhausted classification is configured. */
@@ -489,6 +540,19 @@ public record FleetConfig(
return mcpUrl != null && !mcpUrl.isBlank();
}
/**
* True when the IDE Index MCP should be mounted into a spawned worker (CB-634), pinned to
* the worker's own worktree via the charter. Opt-in per profile, default off.
*/
public boolean hasIdeMcp() {
return ideMcpUrl != null && !ideMcpUrl.isBlank();
}
/** True when this profile mounts any MCP server into its member — the bridge, the IDE, or both. */
public boolean mountsAnyMcp() {
return hasMcp() || hasIdeMcp();
}
/**
* Render this member's tab label (CB-557): {@code {role}}, {@code {profile}},
* {@code {model}} and {@code {n}} are substituted.
@@ -605,6 +669,82 @@ public record FleetConfig(
}
}
/**
* Shared cross-host leader coordination broker: a {@link LeadMailbox} lets two leads on
* different daemons — possibly different hosts — exchange durable messages, which a herdr pane
* injection (how {@code fleet_send} reaches a lead today) cannot do at all. Its mere presence is
* config only in this ticket: nothing here opens a live {@code LeadMailbox} yet, that wiring is
* a separate ticket.
*
* <p>Deliberately a SEPARATE vhost from {@link Broker}, not a reuse of it. {@link Broker} is
* per-fleet — its queues are named by worker session id, and two fleets sharing one broker stay
* isolated by vhost (see {@code Two fleets share one LavinMQ}). Leader coordination is meant to
* cross exactly that boundary: two independently-owned fleets' leads talking to each other. Using
* the same vhost would either leak member traffic across the fleet boundary this is meant to
* cross, or force every fleet's members onto one shared vhost to get leader coordination — a
* second, dedicated vhost keeps "member inboxes stay fleet-local" true while still letting leads
* reach across fleets.
*
* @param uri AMQP connection URI for the coordination vhost, e.g.
* {@code amqp://guest:guest@127.0.0.1:5672/coord}. Blank/{@code null} ⇒ the
* coordinator block is treated as unconfigured. Ignored when {@code uriEnv} is set.
* @param uriEnv name of a host env var holding the AMQP URI, same convention as
* {@link Broker#uriEnv()} — keeps the credential out of the config file. Wins
* over {@code uri} whenever set. Blank/{@code null} ⇒ ignored.
* @param selfId this daemon's own lead coord-id — the name its {@code LeadMailbox} is owned
* under ({@code lead.<selfId>.inbox}), e.g. {@code "mac-opus"}. Blank/
* {@code null} ⇒ kept as {@code null} (no self id configured).
* @param prefetch the consumer's {@code basicQos} prefetch count. {@code null}/non-positive ⇒
* {@link LeadMailbox#DEFAULT_PREFETCH}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Coordinator(String uri, String uriEnv, String selfId, Integer prefetch) {
public Coordinator {
selfId = (selfId == null || selfId.isBlank()) ? null : selfId;
}
/** True when a {@code uriEnv} is configured by name, whether or not its variable resolves. */
public boolean hasUriEnv() {
return uriEnv != null && !uriEnv.isBlank();
}
/**
* True when a usable coordination broker URI is configured (an empty block does not enable
* it). Honors {@code uriEnv} first: if it names a variable that is unset or blank, the
* coordinator is <em>not</em> configured — a bare {@code uri} is only consulted when no
* {@code uriEnv} is set.
*/
public boolean isConfigured() {
return effectiveUri() != null;
}
/**
* The effective AMQP URI to connect with. {@code uriEnv} wins when set (both over
* {@code uri} and alone). When {@code uriEnv} names a variable that is unset or blank,
* returns {@code null} rather than falling back to {@code uri} — an operator who moved to
* the secret store must not silently drop back onto a stale clear-text URI. Returns the
* literal {@code uri} when no {@code uriEnv} is configured.
*/
public String effectiveUri() {
return effectiveUri(System.getenv());
}
/** As {@link #effectiveUri()}, reading the variable value from {@code env} (the injection seam). */
public String effectiveUri(Map<String, String> env) {
if (hasUriEnv()) {
String value = env.get(uriEnv);
return (value != null && !value.isBlank()) ? value : null;
}
return (uri != null && !uri.isBlank()) ? uri : null;
}
/** The prefetch to use, defaulting to {@link LeadMailbox#DEFAULT_PREFETCH} when unset. */
public int prefetchOrDefault() {
return (prefetch != null && prefetch > 0) ? prefetch : LeadMailbox.DEFAULT_PREFETCH;
}
}
/**
* Optional pinned primary terminal config (CB-307). When present with a non-blank
* {@code terminal}, the bridge uses this as the primary's herdr identity instead of
@@ -1139,7 +1279,7 @@ public record FleetConfig(
"bind", "herdrSocket", "profiles", "guard", "worktreeRoot",
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
"memberCredentials");
"memberCredentials", "coordinator");
/** Load and validate config from {@code path}. */
public static FleetConfig load(Path path) {
@@ -1708,9 +1848,11 @@ public record FleetConfig(
// and an upgrade must not change what a running deployment's members inherit.
MemberCredentials mc = memberCredentials != null ? memberCredentials
: new MemberCredentials(null, List.of(), List.of());
// coordinator is left as-is, like broker/primary above: null keeps no LeadMailbox opened,
// and this ticket's Coordinator is config-only anyway (nothing yet reads it at startup).
return new FleetConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
quarantineCooldown, mc);
quarantineCooldown, mc, coordinator);
}
/**
@@ -212,6 +212,17 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
guard.assertWorker(baseUrl); // hard stop before we spawn anything
}
// CB-634: the IDE guidance is delivered as an on-disk CLAUDE.local.md overlay, NOT through
// the charter — the charter returns to role -> reply only. Best-effort: a failed overlay
// must never fail the spawn, and `writeIdeOverlay` no-ops unless the cwd is a provisioned
// worktree (see its .git-file safety gate). The overlay pins, and the auto-open opens, the
// module dir (this repo's pom is in `bridged/`, not at the worktree root) — see ideProjectPath.
if (cfg.hasIdeMcp()) {
String projectPath = PeerLauncher.ideProjectPath(spec.cwd(), cfg.ideProjectDir());
writeIdeOverlay(spec.cwd(), projectPath);
PeerLauncher.openInIde(projectPath, cfg.ideOpenCommand(), log);
}
Map<String, String> workerEnv = baseEnv(cfg);
if (onSubscription) {
// CB-542 belt-and-braces: on the subscription path no guard vets these two keys, and the
@@ -292,26 +303,31 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
*/
private List<String> argvWithFleet(FleetConfig.Profile cfg, LaunchSpec spec) {
String roleCharter = nonBlank(spec.roleCharter());
// CB-634: the IDE guidance is delivered as an on-disk overlay (writeIdeOverlay), not through
// the charter. The charter file is role -> reply only.
String replyCharter = nonBlank(spec.replyCharter());
Path agentFile = agentDefinitionFile(spec.cwd(), spec.role(), ".claude", "agents");
if (!cfg.hasMcp() && roleCharter == null && replyCharter == null && agentFile == null) {
if (!cfg.mountsAnyMcp() && roleCharter == null
&& replyCharter == null && agentFile == null) {
return cfg.argv();
}
List<String> argv = mutableArgv(cfg.argv());
if (cfg.hasMcp()) {
String mcpJson = "{\"mcpServers\":{\"" + PeerLauncher.MCP_MOUNT_NAME
+ "\":{\"type\":\"http\",\"url\":\""
+ cfg.mcpUrl() + "\"}}}";
if (cfg.mountsAnyMcp()) {
argv.add("--mcp-config");
argv.add(mcpJson);
argv.add(mcpConfigJson(cfg));
}
if (roleCharter != null) {
String combined = replyCharter == null ? roleCharter : roleCharter + "\n\n" + replyCharter;
argv.add("--append-system-prompt-file");
argv.add(writeCharterFile(combined).toString());
} else if (replyCharter != null) {
// Combine the charters in order role -> reply, dropping any that are absent. When two
// or more survive they must ride one --append-system-prompt-file (CB-618 forbids the inline
// flag and the file flag together). A lone reply charter keeps its proven inline delivery.
List<String> charters = new java.util.ArrayList<>(2);
if (roleCharter != null) charters.add(roleCharter);
if (replyCharter != null) charters.add(replyCharter);
if (charters.size() == 1 && replyCharter != null && roleCharter == null) {
argv.add("--append-system-prompt");
argv.add(replyCharter);
} else if (!charters.isEmpty()) {
argv.add("--append-system-prompt-file");
argv.add(writeCharterFile(String.join("\n\n", charters)).toString());
}
if (agentFile != null) {
argv.add("--agent");
@@ -320,6 +336,83 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
return argv;
}
/**
* The {@code --mcp-config} JSON for this member: always the bridge mount when {@link
* FleetConfig.Profile#hasMcp()}, plus the IDE Index MCP as a second server named {@code
* intellij} when {@link FleetConfig.Profile#hasIdeMcp()} (CB-634). At least one is present —
* the caller only reaches here when {@link FleetConfig.Profile#mountsAnyMcp()} is true.
*/
private static String mcpConfigJson(FleetConfig.Profile cfg) {
StringBuilder servers = new StringBuilder();
if (cfg.hasMcp()) {
servers.append('"').append(PeerLauncher.MCP_MOUNT_NAME)
.append("\":{\"type\":\"http\",\"url\":\"").append(cfg.mcpUrl()).append("\"}");
}
if (cfg.hasIdeMcp()) {
if (servers.length() > 0) servers.append(',');
servers.append("\"intellij\":{\"type\":\"http\",\"url\":\"")
.append(cfg.ideMcpUrl()).append("\"}");
}
return "{\"mcpServers\":{" + servers + "}}";
}
/**
* Deliver the shared IDE guidance ({@link PeerLauncher#ideOverlayText}) as an on-disk
* {@code CLAUDE.local.md} overlay beside the project's own {@code CLAUDE.md} (CB-634), and
* register the overlay in the repository's common {@code info/exclude} so it never shows as
* untracked (git reads a worktree's excludes from the common dir, not the per-worktree gitdir).
*
* <p><strong>Safety gate:</strong> the overlay is written ONLY when {@code cwd/.git} is a
* <em>regular file</em> — a provisioned worktree keeps a {@code .git} FILE holding a
* {@code gitdir: <path>} pointer, while the primary's real checkout has a {@code .git}
* DIRECTORY. Returning without writing when {@code .git} is a directory is the whole safety of
* the feature: it must never write into a non-worktree cwd, i.e. never clobber a project that
* does not want the overlay.
*
* <p>Best-effort: a failure is logged at debug and swallowed — a failed overlay must never fail
* the spawn.
*
* @param cwd the member's worktree root, where the {@code CLAUDE.local.md} file is written
* @param projectPath the module dir the overlay pins {@code project_path} to (see
* {@link PeerLauncher#ideProjectPath}); equals {@code cwd} when no module subdir
*/
private static void writeIdeOverlay(String cwd, String projectPath) {
try {
Path dotGit = Path.of(cwd, ".git");
if (!Files.isRegularFile(dotGit)) {
// Not a provisioned worktree (primary's real checkout has a .git directory, or the
// cwd is not a repo at all). Never write into it.
return;
}
// The overlay FILE lives at the worktree root (claude-code's cwd), but its CONTENT pins
// project_path to the module dir the IDE opened (projectPath), not the worktree root.
Files.writeString(Path.of(cwd, "CLAUDE.local.md"), PeerLauncher.ideOverlayText(projectPath));
String gitdirLine = Files.readString(dotGit).trim();
Path gitDir = Path.of(gitdirLine.replaceFirst("^gitdir:\\s*", ""));
if (!gitDir.isAbsolute()) {
gitDir = Path.of(cwd).resolve(gitDir).normalize();
}
// git reads info/exclude from the COMMON dir, never the per-worktree gitdir (only
// info/sparse-checkout is per-worktree). A provisioned worktree's gitdir is
// <common>/worktrees/<name>, so the common dir is two levels up; writing the entry into
// the per-worktree gitdir leaves it un-honoured and the overlay shows as untracked.
Path commonDir = gitDir;
if (gitDir.getParent() != null && gitDir.getParent().getFileName() != null
&& "worktrees".equals(gitDir.getParent().getFileName().toString())) {
commonDir = gitDir.getParent().getParent();
}
Path exclude = commonDir.resolve("info").resolve("exclude");
Files.createDirectories(exclude.getParent());
String overlayLine = "CLAUDE.local.md";
if (!Files.exists(exclude) || Files.readAllLines(exclude).stream().noneMatch(overlayLine::equals)) {
Files.writeString(exclude, (Files.exists(exclude) ? System.lineSeparator() : "")
+ overlayLine + System.lineSeparator());
}
} catch (Exception e) {
log.debug("cannot write IDE overlay into worktree '{}'", cwd, e);
}
}
/** {@code s}, or {@code null} when {@code s} is null/blank — the charter-presence test used above. */
private static String nonBlank(String s) {
return (s == null || s.isBlank()) ? null : s;
@@ -24,6 +24,9 @@ import java.util.function.Function;
import java.util.function.LongSupplier;
import java.util.function.Supplier;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* The {@link HerdrPeerLauncher} adapter for <strong>opencode</strong> — an open-source,
* provider-agnostic terminal coding agent. Its whole reason for existing is to prove the
@@ -50,6 +53,8 @@ import java.util.function.Supplier;
*/
public final class OpenCodeLauncher extends HerdrPeerLauncher {
private static final Logger log = LoggerFactory.getLogger(OpenCodeLauncher.class);
/** Label prefix for this adapter's herdr agent names (drives naming + orphan reap). */
private static final String NAME_PREFIX = "opencode";
@@ -204,9 +209,10 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
@Override
protected Launch buildLaunch(FleetConfig.Profile cfg, LaunchSpec spec) {
Map<String, String> workerEnv = baseEnv(cfg);
// A config file is needed for the bridge MCP mount, a member charter, or a pinned endpoint (CB-508).
if (cfg.hasMcp() || spec.charter() != null || hasCustomProvider(cfg)) {
workerEnv.put("OPENCODE_CONFIG", writeConfig(cfg, spec.charter()).toString());
// A config file is needed for the bridge MCP mount, a member charter, the IDE MCP (+ its
// guidance overlay, CB-634), or a pinned endpoint (CB-508).
if (cfg.hasMcp() || cfg.hasIdeMcp() || spec.charter() != null || hasCustomProvider(cfg)) {
workerEnv.put("OPENCODE_CONFIG", writeConfig(cfg, spec.charter(), spec.cwd()).toString());
}
applyGitToken(workerEnv, cfg);
List<String> argv = argvWithResume(argvWithModel(argvWithAuto(cfg), cfg), spec.resumeSessionId());
@@ -290,7 +296,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* {@code OPENCODE_CONFIG}. The dir is unique per spawn so concurrent workers never race on it;
* it is best-effort cleaned on JVM exit (worker config is disposable — regenerated every spawn).
*/
private Path writeConfig(FleetConfig.Profile cfg, String charterText) {
private Path writeConfig(FleetConfig.Profile cfg, String charterText, String cwd) {
try {
Path dir = Files.createTempDirectory(configRoot, "bridged-opencode-");
dir.toFile().deleteOnExit();
@@ -321,11 +327,35 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
root.putArray("instructions").add(charter.toAbsolutePath().toString());
}
if (cfg.hasMcp()) {
ObjectNode mount = root.putObject("mcp").putObject(PeerLauncher.MCP_MOUNT_NAME);
mount.put("type", "remote");
mount.put("url", cfg.mcpUrl());
mount.put("enabled", true);
if (cfg.hasMcp() || cfg.hasIdeMcp()) {
// One shared mcp node for both servers — putObject would replace the node (and thus
// the other server) on the second call, so build into a single get-or-create node.
ObjectNode mcp = root.withObject("mcp");
if (cfg.hasMcp()) {
ObjectNode mount = mcp.putObject(PeerLauncher.MCP_MOUNT_NAME);
mount.put("type", "remote");
mount.put("url", cfg.mcpUrl());
mount.put("enabled", true);
}
// CB-634: mount the IDE Index MCP in the same shape as the bridge remote server, and
// deliver its guidance via the instructions array (opencode does not read
// CLAUDE.local.md) rather than any system-prompt string.
if (cfg.hasIdeMcp()) {
ObjectNode ide = mcp.putObject("intellij");
ide.put("type", "remote");
ide.put("url", cfg.ideMcpUrl());
ide.put("enabled", true);
// CB-634: pin the overlay and open the IDE at the module dir (this repo's pom is
// in `bridged/`, not at the worktree root) — see PeerLauncher.ideProjectPath.
String projectPath = PeerLauncher.ideProjectPath(cwd, cfg.ideProjectDir());
Path rules = dir.resolve("ide-rules.md");
Files.writeString(rules, PeerLauncher.ideOverlayText(projectPath));
rules.toFile().deleteOnExit();
// The array may already hold the member-charter path; withArray gets-or-creates.
root.withArray("instructions").add(rules.toAbsolutePath().toString());
PeerLauncher.openInIde(projectPath, cfg.ideOpenCommand(), log);
}
}
if (hasCustomProvider(cfg)) {
addCustomProvider(root, cfg);
@@ -0,0 +1,422 @@
package dev.ltms.fleet.msg;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.Connection;
import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import com.rabbitmq.client.Recoverable;
import com.rabbitmq.client.RecoveryListener;
import com.rabbitmq.client.Return;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.NavigableMap;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentSkipListMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
/**
* Durable, AMQP-backed mailbox for lead-to-lead messages across daemons — including daemons on
* different hosts, where a herdr pane injection (how {@code fleet_send} reaches a lead today)
* cannot reach at all. The broker is the only medium two independently-owned daemons share, which
* is exactly why {@link AmqpReplyInbox}'s javadoc already calls out "one gateway may publish to an
* agent owned by another gateway" (CB-308 federation) as the reason {@code publish} and
* {@code own}/consume are separate operations there — this class leans on the same split.
*
* <p><strong>Single-target, unlike {@link AmqpReplyInbox}.</strong> {@code AmqpReplyInbox}
* multiplexes many workers' reply queues under one gateway connection. A {@code LeadMailbox}
* instance is simpler: it owns exactly <em>one</em> queue — this daemon's own
* {@code lead.<selfCoordId>.inbox} — declared and consumed the moment it is constructed. There is
* no {@code own}/{@code release} pair to call separately; a daemon either runs a {@code LeadMailbox}
* for its own coord-id, or it does not run one at all.
*
* <p><strong>Consume-and-hold with deferred manual ack</strong> — same mapping as
* {@code AmqpReplyInbox}. The constructor declares the durable queue and starts a manual-ack
* consumer that pulls persistent messages into an in-memory {@code held} map (keyed by
* {@link LeadMessage#msgId()}) but does not ack them. {@link #peek} returns a non-destructive
* snapshot; {@link #ack} acks the broker delivery-tag and drops the entry. A message that is never
* acked (a crash, a bounce) survives — the broker redelivers it to the next connection that owns
* the queue.
*
* <p><strong>Publishing does not imply owning.</strong> {@link #publish} sends to
* {@code lead.<toCoordId>.inbox} over a dedicated confirm-mode channel; it never declares that
* queue as owned and never attaches a consumer to it. A sender that has never opened its own
* {@code LeadMailbox} for {@code toCoordId} can still publish to it, exactly as CB-308 federation
* requires. Publish blocks for the broker's publisher confirm (persistent delivery, {@code
* mandatory=true}) and throws {@link IllegalStateException} on an unroutable return, a nack, or a
* timeout — the caller must not report success for a black-holed message.
*
* <p><strong>Recovery.</strong> The connection is opened with automatic + topology recovery
* enabled, mirroring {@code AmqpReplyInbox}: on reconnect the broker hands out fresh delivery tags,
* so the held snapshot is cleared (dedup by {@code msgId} still prevents any double-queue on
* redelivery) and any publish still awaiting its confirm is failed rather than left to idle out
* the confirm timeout against a sequence number that means nothing on the new channel.
*/
public final class LeadMailbox implements AutoCloseable {
private static final Logger log = LoggerFactory.getLogger(LeadMailbox.class);
private static final String QUEUE_PREFIX = "lead.";
private static final String QUEUE_SUFFIX = ".inbox";
/** The prefetch used when a caller does not pass an explicit value to {@link #open(String, String, int)}. */
public static final int DEFAULT_PREFETCH = 32;
/** How long {@link #publish} waits for its publisher confirm before failing the call. */
private static final long CONFIRM_TIMEOUT_MS = 10_000L;
private static final ObjectMapper MAPPER = new ObjectMapper();
private final Connection connection;
private final String selfCoordId;
private final Channel channel;
/** All consume-channel operations (declare/consume/ack) serialize on this — a Channel is not thread-safe. */
private final Object channelLock = new Object();
/** msgId → held delivery, for this mailbox's own queue only (there is exactly one). */
private final LinkedHashMap<String, Held> held = new LinkedHashMap<>();
/**
* A dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume + ack)
* so a publish confirm round trip never blocks under {@link #channelLock} and stalls an ack.
*/
private final Channel publishChannel;
private final Object publishChannelLock = new Object();
/** In-flight publishes awaiting their confirm, keyed by the publish channel's sequence number. */
private final ConcurrentSkipListMap<Long, Pending> pendingBySeq = new ConcurrentSkipListMap<>();
/** The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery tag. */
private final ConcurrentHashMap<String, Pending> pendingByMsgId = new ConcurrentHashMap<>();
/** A message pulled off the broker but not yet acked: its delivery-tag plus the deserialized envelope. */
private record Held(long deliveryTag, LeadMessage message) {}
/** A publish awaiting its confirm; {@link #returned} records whether the broker already returned it. */
private static final class Pending {
final String msgId;
final CompletableFuture<Void> confirmed = new CompletableFuture<>();
volatile boolean returned;
Pending(String msgId) {
this.msgId = msgId;
}
}
/**
* Connect to {@code uri} (the shared cross-host coordination vhost, e.g.
* {@code amqp://guest:guest@127.0.0.1:5672/coord}) and own {@code selfCoordId}'s mailbox, with
* {@link #DEFAULT_PREFETCH}.
*/
public static LeadMailbox open(String uri, String selfCoordId) {
return open(uri, selfCoordId, DEFAULT_PREFETCH);
}
/** As {@link #open(String, String)}, with an explicit consumer prefetch. */
public static LeadMailbox open(String uri, String selfCoordId, int prefetch) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares the queue and re-attaches the consumer.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
return new LeadMailbox(factory.newConnection("bridged-lead-mailbox"), selfCoordId, prefetch);
} catch (Exception e) {
throw new IllegalStateException("cannot connect to AMQP coordination broker at " + uri, e);
}
}
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for tests). */
LeadMailbox(Connection connection, String selfCoordId) {
this(connection, selfCoordId, DEFAULT_PREFETCH);
}
/** As above, with an explicit prefetch (injection seam for tests). */
LeadMailbox(Connection connection, String selfCoordId, int prefetch) {
this.connection = connection;
this.selfCoordId = selfCoordId;
try {
this.channel = connection.createChannel();
// Bound the held backlog — must be set before basicConsume.
this.channel.basicQos(prefetch);
this.publishChannel = connection.createChannel();
this.publishChannel.confirmSelect();
this.publishChannel.addReturnListener(this::onReturn);
this.publishChannel.addConfirmListener(this::onAck, this::onNack);
own();
} catch (IOException e) {
throw new IllegalStateException("cannot open AMQP channel", e);
}
// On automatic recovery the broker redelivers unacked messages with FRESH delivery-tags; the
// tags we were holding are now stale. Drop the held snapshot so the re-attached consumer
// repopulates it with valid tags (dedup by msgId still prevents any double-queue). Any publish
// confirm still in flight when the connection dropped is equally stale — fail it now rather
// than let it silently ride out CONFIRM_TIMEOUT_MS.
if (connection instanceof Recoverable recoverable) {
recoverable.addRecoveryListener(new RecoveryListener() {
@Override
public void handleRecovery(Recoverable recoverable) {
synchronized (held) {
held.clear();
}
failPendingPublishesOnRecovery();
log.info("AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery");
}
@Override
public void handleRecoveryStarted(Recoverable recoverable) {
// no-op: we act once recovery completes
}
});
}
}
/** Declare + consume this daemon's own {@code lead.<selfCoordId>.inbox}. Called once, at construction. */
private void own() throws IOException {
String queue = queueName(selfCoordId);
synchronized (channelLock) {
channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
channel.basicConsume(queue, false, deliverCallback(), _ -> { }); // autoAck=false: manual ack
}
log.debug("lead mailbox owns queue {} for coord-id {}", queue, selfCoordId);
}
/**
* Publish {@code msg} to {@code toCoordId}'s mailbox and block until the broker's publisher
* confirm for it lands. Does <em>not</em> imply owning or consuming {@code toCoordId}'s queue.
* Throws {@link IllegalStateException} if the message is returned as unroutable, nacked, or not
* confirmed within {@link #CONFIRM_TIMEOUT_MS} — the caller must treat that as a failed publish,
* not a lost-and-forgotten one.
*/
public void publish(String toCoordId, LeadMessage msg) {
byte[] body;
try {
body = MAPPER.writeValueAsBytes(msg);
} catch (JsonProcessingException e) {
throw new IllegalStateException("cannot serialize lead message " + msg.msgId(), e);
}
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
.messageId(msg.msgId())
.deliveryMode(2) // persistent — survives a broker restart
.contentType("application/json")
.build();
Pending pending = new Pending(msg.msgId());
long seq;
synchronized (publishChannelLock) {
seq = publishChannel.getNextPublishSeqNo();
pendingBySeq.put(seq, pending);
pendingByMsgId.put(msg.msgId(), pending);
try {
publishChannel.basicPublish("", queueName(toCoordId), true, props, body);
} catch (IOException e) {
pendingBySeq.remove(seq, pending);
pendingByMsgId.remove(msg.msgId(), pending);
throw new IllegalStateException("cannot publish lead message to " + queueName(toCoordId), e);
}
}
try {
pending.confirmed.get(CONFIRM_TIMEOUT_MS, TimeUnit.MILLISECONDS);
} catch (ExecutionException e) {
Throwable cause = e.getCause();
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
} catch (TimeoutException e) {
throw new IllegalStateException("publish confirm for lead message " + msg.msgId() + " to "
+ queueName(toCoordId) + " timed out after " + CONFIRM_TIMEOUT_MS
+ "ms — broker may be unreachable or overloaded", e);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new IllegalStateException("interrupted awaiting publish confirm for " + msg.msgId(), e);
} finally {
pendingBySeq.remove(seq, pending);
pendingByMsgId.remove(msg.msgId(), pending);
}
}
/** Non-destructive FIFO snapshot of this mailbox's currently-held messages. */
public List<LeadMessage> peek() {
synchronized (held) {
return held.values().stream().map(Held::message).toList();
}
}
/** Convenience: {@link #peek} the current snapshot, then {@link #ack} every message in it. */
public List<LeadMessage> drain() {
List<LeadMessage> snapshot = peek();
snapshot.forEach(m -> ack(m.msgId()));
return snapshot;
}
/** Remove the held message {@code msgId} and ack it on the broker. No-op if not held. */
public void ack(String msgId) {
Held h;
synchronized (held) {
h = held.remove(msgId);
}
if (h == null) {
return; // never held (or already acked) — no-op
}
try {
synchronized (channelLock) {
channel.basicAck(h.deliveryTag(), false);
}
} catch (IOException e) {
// Ack didn't reach the broker: restore the entry so a later ack (or a redelivery after
// reconnect) can retry. Keeps the at-least-once contract — a message is never silently lost.
synchronized (held) {
held.putIfAbsent(msgId, h);
}
throw new IllegalStateException("cannot ack lead message " + msgId, e);
}
}
private DeliverCallback deliverCallback() {
return (_, delivery) -> {
long tag = delivery.getEnvelope().getDeliveryTag();
LeadMessage msg;
try {
msg = MAPPER.readValue(delivery.getBody(), LeadMessage.class);
} catch (IOException e) {
// A malformed body can never be dedup-keyed or handed to a caller; ack it so the
// broker does not redeliver it forever, and log loudly since this should never happen
// for a producer that only ever calls publish(String, LeadMessage).
log.warn("dropping malformed lead-mailbox delivery (tag {}): {}", tag, e.toString());
synchronized (channelLock) {
channel.basicAck(tag, false);
}
return;
}
boolean duplicate;
synchronized (held) {
if (held.containsKey(msg.msgId())) {
duplicate = true;
} else {
held.put(msg.msgId(), new Held(tag, msg));
duplicate = false;
}
}
if (duplicate) {
// Redelivered duplicate: ack the new tag and drop it so the broker stops resending.
synchronized (channelLock) {
channel.basicAck(tag, false);
}
}
};
}
/** Broker return for an unroutable {@code mandatory} publish — arrives BEFORE its confirm. */
private void onReturn(Return r) {
String msgId = r.getProperties() == null ? null : r.getProperties().getMessageId();
Pending pending = msgId == null ? null : pendingByMsgId.get(msgId);
if (pending != null) {
pending.returned = true;
} else {
log.warn("AMQP return for lead message {} (routingKey={}, {} {}) with no matching in-flight publish"
+ " — already resolved by a prior confirm", msgId, r.getRoutingKey(), r.getReplyCode(),
r.getReplyText());
}
}
private void onAck(long seq, boolean multiple) {
resolveConfirm(seq, multiple, true);
}
private void onNack(long seq, boolean multiple) {
resolveConfirm(seq, multiple, false);
}
/**
* Resolve every pending publish covered by this confirm (a single seq, or — {@code multiple} —
* every seq up to and including it). Checks {@link Pending#returned} at confirm time: since the
* broker's return for an unroutable message always precedes its confirm, an ack that arrives after
* a return means "confirmed but never routed", not "durably queued".
*/
private void resolveConfirm(long seq, boolean multiple, boolean ack) {
NavigableMap<Long, Pending> covered = multiple
? pendingBySeq.headMap(seq, true)
: pendingBySeq.subMap(seq, true, seq, true);
for (var it = covered.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
if (ack && !pending.returned) {
pending.confirmed.complete(null);
} else if (ack) {
pending.confirmed.completeExceptionally(new IllegalStateException(
"lead message " + pending.msgId + " was returned as unroutable (mailbox not owned)"));
} else {
pending.confirmed.completeExceptionally(new IllegalStateException(
"broker nacked publish of lead message " + pending.msgId));
}
}
}
/**
* Fail every publish still awaiting its confirm — their sequence numbers are stale after
* recovery. Guarded by {@link #publishChannelLock}, the same lock {@link #publish} holds while
* it takes its sequence number and registers its {@link Pending} — see
* {@code AmqpReplyInbox.failPendingPublishesOnRecovery}'s javadoc for the full race analysis this
* mirrors. Package-private only so a unit test can drive it directly without a live broker
* reconnect.
*/
void failPendingPublishesOnRecovery() {
synchronized (publishChannelLock) {
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
pending.confirmed.completeExceptionally(new IllegalStateException(
"AMQP connection recovered mid-publish; confirm status of lead message "
+ pending.msgId + " is unknown"));
}
}
}
/**
* Fail every publish still awaiting its confirm with a clear, immediate error instead of leaving
* it to time out after {@link #CONFIRM_TIMEOUT_MS} once the channels are closed underneath it.
*/
private void failPendingPublishesOnClose() {
synchronized (publishChannelLock) {
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
pending.confirmed.completeExceptionally(new IllegalStateException(
"lead mailbox closed while publish of lead message " + pending.msgId
+ " was still awaiting its confirm"));
}
}
}
/** The durable queue name a coord-id's mailbox lives on: {@code lead.<coordId>.inbox}. */
public static String queueName(String coordId) {
return QUEUE_PREFIX + coordId + QUEUE_SUFFIX;
}
@Override
public void close() {
failPendingPublishesOnClose();
try {
channel.close();
} catch (Exception e) {
log.debug("AMQP lead mailbox channel close: {}", e.toString());
}
try {
publishChannel.close();
} catch (Exception e) {
log.debug("AMQP lead mailbox publish channel close: {}", e.toString());
}
try {
connection.close();
} catch (Exception e) {
log.debug("AMQP lead mailbox connection close: {}", e.toString());
}
}
}
@@ -0,0 +1,22 @@
package dev.ltms.fleet.msg;
/**
* Wire envelope for a lead-to-lead message carried over {@link LeadMailbox}.
*
* <p>Unlike {@link ReplyInbox.InboxMessage} (a worker→primary reply, addressed only by the single
* gateway that owns the worker), a lead message crosses independently-owned daemons — possibly on
* different hosts — so it carries an explicit sender ({@code from}) as well as the recipient
* ({@code to}): the recipient needs the sender's coord-id to reply back.
*
* <p>{@code from} and {@code to} are globally-unique lead coordination ids (e.g. {@code "mac-opus"},
* {@code "fleet01-lead"}) — NOT herdr terminal ids. A herdr terminal id is meaningful only on the
* host that owns it, so it cannot address a lead running on another daemon; a coord-id is chosen
* by configuration ({@code coordinator.selfId}) precisely so it means the same thing everywhere.
*
* @param msgId idempotency id; a redelivered duplicate (at-least-once delivery) is deduped on this
* @param from the sending lead's coord-id
* @param to the receiving lead's coord-id — identifies the mailbox this message is held on
* @param content the message text
*/
public record LeadMessage(String msgId, String from, String to, String content) {
}
@@ -1,8 +1,11 @@
package dev.ltms.fleet.peer;
import java.nio.file.Path;
import java.util.List;
import java.util.Set;
import org.slf4j.Logger;
/**
* SPI for materializing a connected peer — the only way the bridge core creates or tears down
* a peer process. Every launcher is a first-party, in-tree adapter selected by (future) profile
@@ -30,6 +33,79 @@ public interface PeerLauncher {
*/
String MCP_MOUNT_NAME = "fleet";
/**
* Shared IDE-guidance text (CB-634), delivered per-backend as an on-disk overlay rather than
* any one adapter's system-prompt charter, so a project's own {@code CLAUDE.md} is never
* clobbered. It pins every {@code ide_*} call to the member's own worktree, which is the whole
* point of the mechanism. Both launchers render their own overlay from this single source.
*
* @param projectPath the path the member must pin every {@code ide_*} call to — the module dir
* IntelliJ opened as the project, which is {@link #ideProjectPath} of the
* member's own worktree (the worktree root when no module subdir is set)
*/
static String ideOverlayText(String projectPath) {
return "## IDE code intelligence — your worktree only\n"
+ "An IntelliJ IDE Index MCP server is mounted as `mcp__intellij__ide_*`. Prefer it "
+ "over `grep`/`find` for symbol lookups, references, call and type hierarchy, and "
+ "diagnostics — it resolves the real AST, text search does not.\n\n"
+ "Every `ide_*` call MUST pass `project_path: \"" + projectPath + "\"` — your own "
+ "worktree — and never any other path. A call without it errors "
+ "`multiple_projects_open`; a call with a different path reads another checkout, "
+ "not your changes. This is not the primary's IDE: it is your worktree, pinned to "
+ "you.";
}
/**
* The absolute path IntelliJ must open as the project, and the {@code project_path} the overlay
* pins (CB-634). It is {@code cwd} resolved against {@code ideProjectDir}. The distinction
* matters because this repo (like {@code fleet/fleetd}) keeps its Maven module in a subdir
* ({@code bridged/}), not at the worktree root: opening the root imports no module and
* {@code ide_*} resolves nothing, so the module dir is the correct pin and open target.
*
* @param cwd the member's worktree root
* @param ideProjectDir repo-relative module subdir, or {@code null}/blank for the worktree root
* @return the absolute, normalized module dir as a string
*/
static String ideProjectPath(String cwd, String ideProjectDir) {
Path base = Path.of(cwd);
if (ideProjectDir == null || ideProjectDir.isBlank()) {
return base.toString();
}
return base.resolve(ideProjectDir).normalize().toString();
}
/**
* Best-effort: open {@code projectPath} in the host IDE by running {@code openCommand} with
* every {@code {dir}} replaced by {@code projectPath} (CB-634 auto-open). The command runs
* through {@code /bin/sh -c} so an operator can set env inline — e.g.
* {@code "env DISPLAY=:10.0 idea {dir}"} — because the daemon's own env may lack {@code DISPLAY}.
*
* <p>A blank command is a no-op: the profile opted into IDE MCP but not auto-open, so the
* operator opens the module by hand. The child process is detached and its exit is not awaited;
* any failure is logged and swallowed, because a member must spawn whether or not an IDE is
* running. There is no close half yet (CB-634 defers it): an opened module stays open until the
* operator closes it, and opening the same module again just refocuses it.
*
* @param projectPath the module dir to open (typically {@link #ideProjectPath})
* @param openCommand the host command template, with {@code {dir}} substituted; null/blank ⇒ no-op
* @param log the calling launcher's logger, for the best-effort WARN
*/
static void openInIde(String projectPath, String openCommand, Logger log) {
if (openCommand == null || openCommand.isBlank()) {
return;
}
String cmd = openCommand.replace("{dir}", projectPath);
try {
new ProcessBuilder("/bin/sh", "-c", cmd)
.redirectOutput(ProcessBuilder.Redirect.DISCARD)
.redirectError(ProcessBuilder.Redirect.DISCARD)
.start();
log.info("CB-634 auto-open: launched IDE open for {}", projectPath);
} catch (Exception e) {
log.warn("CB-634 auto-open of '{}' failed (member still spawns): {}", projectPath, e.getMessage());
}
}
/**
* The set of {@link Capability capabilities} this launcher declares. A peer whose profile
* opts into a git-forge token should include {@link Capability#SELF_PR}; the base set for
@@ -1,6 +1,7 @@
package dev.ltms.fleet.config;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.msg.LeadMailbox;
import dev.ltms.fleet.peer.MemberRole;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
@@ -932,6 +933,66 @@ class FleetConfigTest {
assertFalse(cfg.broker().isConfigured(), "an empty uri must not enable AMQP");
}
@Test
void absentCoordinatorBlockLeavesCoordinatorNull(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-coordinator.yaml");
Files.writeString(f, "bind:\n port: 8080\n");
FleetConfig cfg = FleetConfig.load(f);
assertNull(cfg.coordinator(), "no coordinator: block → null → no lead mailbox is opened");
}
@Test
void coordinatorBlockParses(@TempDir Path dir) throws Exception {
Path f = dir.resolve("coordinator.yaml");
Files.writeString(f, """
bind:
port: 8080
coordinator:
uri: amqp://guest:guest@127.0.0.1:5672/coord
selfId: mac-opus
prefetch: 16
""");
FleetConfig cfg = FleetConfig.load(f);
assertNotNull(cfg.coordinator());
assertTrue(cfg.coordinator().isConfigured(), "a non-blank uri enables the coordinator");
assertEquals("amqp://guest:guest@127.0.0.1:5672/coord", cfg.coordinator().uri());
assertEquals("mac-opus", cfg.coordinator().selfId());
assertEquals(16, cfg.coordinator().prefetchOrDefault());
}
@Test
void coordinatorBlockWithBlankUriStaysUnconfigured(@TempDir Path dir) throws Exception {
Path f = dir.resolve("coordinator-blank.yaml");
Files.writeString(f, "bind:\n port: 8080\ncoordinator:\n uri: \"\"\n");
FleetConfig cfg = FleetConfig.load(f);
assertNotNull(cfg.coordinator());
assertFalse(cfg.coordinator().isConfigured(), "an empty uri must not enable the coordinator");
assertNull(cfg.coordinator().selfId(), "a blank/absent selfId stays null, never coerced to empty");
}
@Test
void coordinatorEffectiveUriHonorsUriEnv() {
FleetConfig.Coordinator withEnv = new FleetConfig.Coordinator(
"amqp://stale-clear-text@127.0.0.1:5672/coord", "LEAD_COORD_URI", "fleet01-lead", null);
assertEquals("amqp://from-env@127.0.0.1:5672/coord",
withEnv.effectiveUri(Map.of("LEAD_COORD_URI", "amqp://from-env@127.0.0.1:5672/coord")),
"uriEnv wins over a literal uri when its variable resolves");
assertNull(withEnv.effectiveUri(Map.of()),
"an unset uriEnv variable must not fall back to the literal uri");
assertNull(withEnv.effectiveUri(Map.of("LEAD_COORD_URI", " ")),
"a blank uriEnv variable must not fall back to the literal uri");
FleetConfig.Coordinator noEnv = new FleetConfig.Coordinator(
"amqp://guest:guest@127.0.0.1:5672/coord", null, null, null);
assertEquals("amqp://guest:guest@127.0.0.1:5672/coord", noEnv.effectiveUri(Map.of()),
"the literal uri is used when no uriEnv is configured");
assertEquals(LeadMailbox.DEFAULT_PREFETCH, noEnv.prefetchOrDefault());
}
@Test
void absentPrimaryBlockLeavesPrimaryNull(@TempDir Path dir) throws Exception {
Path f = dir.resolve("no-primary.yaml");
@@ -444,7 +444,7 @@ class FleetMcpTest {
FleetConfig.Profile wcfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN", null,
"tab", "bridged-workers", "worker: {profile} #{n}", null,
null, null, null, null, null, null, null, 0, null, null, null);
null, null, null, null, null, null, null, 0, null, null, null, null);
Map<String, FleetConfig.Profile> profiles = Map.of(wcfg.profile(), wcfg);
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(
new AgentControl(h), new WorkspaceControl(h), new SubscriptionGuard(Set.of("gx00.gw")),
@@ -12,6 +12,7 @@ import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.peer.SpawnRequest;
import org.junit.jupiter.api.Test;
@@ -63,6 +64,174 @@ class ClaudeCodeLauncherTest {
assertTrue(args.stream().anyMatch(a -> a.contains("fleet_reply")), "reply charter present");
}
// CB-634: a profile with ideMcpUrl set mounts the IDE Index MCP as a second server and pins
// every ide_* call to the member's own worktree via the charter.
/** A profile carrying an ideMcpUrl (plus optional bridge mcpUrl and cwd). ideMcpUrl is the last record component. */
private FleetConfig.Profile ideProfile(String mcpUrl, String ideMcpUrl, String cwd) {
return new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
List.of("claude"), "tab", "bridged-workers", "w #{n}", mcpUrl, cwd, null,
null, null, null, null, null, null, null, null, null, ideMcpUrl);
}
/** As {@link #ideProfile} but carrying the CB-634 auto-open fields (module subdir + open command). */
private FleetConfig.Profile ideProfileModule(String ideMcpUrl, String cwd, String ideProjectDir,
String ideOpenCommand) {
return new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
List.of("claude"), "tab", "bridged-workers", "w #{n}", null, cwd, null,
null, null, null, null, null, null, null, null, null, ideMcpUrl, ideProjectDir, ideOpenCommand);
}
private ClaudeCodeLauncher launcher(FakeHerdr herdr, FleetConfig.Profile cfg) {
return new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
}
@Test
void mountsIdeMcpAsSecondServerWhenIdeMcpUrlSet() {
FakeHerdr herdr = new FakeHerdr();
launcher(herdr, ideProfile("http://127.0.0.1:8765/mcp",
"http://127.0.0.1:29170/index-mcp/streamable-http", null)).spawn();
List<String> args = spawnedArgs(herdr);
assertTrue(args.contains("--mcp-config"));
String json = args.get(args.indexOf("--mcp-config") + 1);
assertTrue(json.contains("\"fleet\"") && json.contains("http://127.0.0.1:8765/mcp"),
"bridge mount still present: " + json);
assertTrue(json.contains("\"intellij\"") && json.contains("29170"),
"IDE Index MCP mounted as a second server named intellij: " + json);
}
@Test
void ideMcpUrlAloneStillEmitsTheMount() {
FakeHerdr herdr = new FakeHerdr();
launcher(herdr, ideProfile(null, "http://127.0.0.1:29170/index-mcp/streamable-http", null)).spawn();
List<String> args = spawnedArgs(herdr);
assertTrue(args.contains("--mcp-config"),
"the mount gate fires on ideMcpUrl alone, not only on mcpUrl");
String json = args.get(args.indexOf("--mcp-config") + 1);
assertTrue(json.contains("\"intellij\""), "IDE server present: " + json);
assertFalse(json.contains("\"fleet\""), "no bridge server when mcpUrl is unset: " + json);
}
@Test
void roleAndReplyOnlyComposeTheCharterFile_NotIdeGuidance() {
FakeHerdr herdr = new FakeHerdr();
String roleCharter = "You review changes.";
String worktree = "/tmp/.fleet-worktrees/rev-1";
FleetConfig.Profile cfg = ideProfile("http://127.0.0.1:8765/mcp",
"http://127.0.0.1:29170/index-mcp/streamable-http", worktree);
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("ltms-local", null, null, null, null, MemberRole.REVIEWER));
List<String> args = spawnedArgs(herdr);
assertTrue(args.stream().noneMatch(a -> a.contains("\n")),
"no argv element may be multi-line: " + args);
assertFalse(args.contains("--append-system-prompt"),
"CB-618: role + reply ride one --append-system-prompt-file, never both flags: " + args);
int fileFlag = args.indexOf("--append-system-prompt-file");
assertTrue(fileFlag >= 0, "the combined charter is mounted via file: " + args);
assertDoesNotThrow(() -> {
String w = Files.readString(Path.of(args.get(fileFlag + 1)));
assertTrue(w.startsWith(roleCharter), "role charter first: " + w);
assertFalse(w.contains("project_path"),
"CB-634: the IDE guidance is delivered as an on-disk overlay, not the charter: " + w);
assertTrue(w.endsWith(HerdrPeerLauncher.REPLY_CHARTER),
"the reply charter is last — it is the rule that must survive: " + w);
}, "the --append-system-prompt-file path must be a readable file");
}
// CB-634: the IDE guidance is delivered as a CLAUDE.local.md overlay (written only into a
// provisioned worktree — cwd with a `.git` FILE) and registered in the repository's COMMON
// info/exclude. git reads a worktree's excludes from the common dir, not the per-worktree
// gitdir (only info/sparse-checkout is per-worktree), so a real provisioned layout
// <common>/worktrees/<name> must land the entry in <common>/info/exclude.
@Test
void writeIdeOverlayWritesClaudeLocalAndAddsItToCommonInfoExclude(@TempDir Path root) throws Exception {
Path worktree = Files.createDirectory(root.resolve("worktree"));
// Real worktree layout: the .git FILE points at <common>/worktrees/<name>.
Path commonDir = Files.createDirectory(root.resolve("dotgit"));
Path gitDir = Files.createDirectories(commonDir.resolve("worktrees").resolve("wt1"));
Files.writeString(worktree.resolve(".git"), "gitdir: " + gitDir);
FakeHerdr herdr = new FakeHerdr();
launcher(herdr, ideProfile(null, "http://127.0.0.1:29170/index-mcp/streamable-http",
worktree.toString())).spawn();
Path overlay = worktree.resolve("CLAUDE.local.md");
assertTrue(Files.exists(overlay), "the overlay is written beside the project's CLAUDE.md");
assertTrue(Files.readString(overlay).contains("project_path: \"" + worktree + "\""),
"the overlay pins every ide_* call to the member's own worktree");
Path commonExclude = commonDir.resolve("info").resolve("exclude");
assertTrue(Files.exists(commonExclude), "info/exclude is created in the COMMON git dir");
assertTrue(Files.readAllLines(commonExclude).contains("CLAUDE.local.md"),
"the overlay is registered in the common exclude git actually honours for a worktree");
assertFalse(Files.exists(gitDir.resolve("info").resolve("exclude")),
"the entry must NOT go to the per-worktree gitdir, which git ignores for excludes");
}
// CB-634 auto-open pin fix: for a repo whose Maven module is a subdir (this repo's pom lives in
// `bridged/`, not at the worktree root), the overlay must pin project_path to the MODULE dir the
// IDE opened, not the worktree root — else ide_* resolves against a root that imports no module.
@Test
void writeIdeOverlayPinsModuleDirWhenIdeProjectDirSet(@TempDir Path root) throws Exception {
Path worktree = Files.createDirectory(root.resolve("worktree"));
Path commonDir = Files.createDirectory(root.resolve("dotgit"));
Path gitDir = Files.createDirectories(commonDir.resolve("worktrees").resolve("wt1"));
Files.writeString(worktree.resolve(".git"), "gitdir: " + gitDir);
FakeHerdr herdr = new FakeHerdr();
// ideProjectDir "bridged" ⇒ the pin is <worktree>/bridged, not <worktree>.
launcher(herdr, ideProfileModule("http://127.0.0.1:29170/index-mcp/streamable-http",
worktree.toString(), "bridged", null)).spawn();
Path overlay = worktree.resolve("CLAUDE.local.md");
assertTrue(Files.exists(overlay), "the overlay file still lives at the worktree root");
String module = worktree.resolve("bridged").toString();
assertTrue(Files.readString(overlay).contains("project_path: \"" + module + "\""),
"the overlay pins the MODULE dir the IDE opened, not the worktree root");
assertFalse(Files.readString(overlay).contains("project_path: \"" + worktree + "\""),
"the bare worktree root must NOT be the pin when a module subdir is set");
}
@Test
void ideProjectPathResolvesModuleSubdirAndDefaultsToWorktreeRoot() {
assertEquals("/wt/x/bridged", PeerLauncher.ideProjectPath("/wt/x", "bridged"),
"a module subdir resolves under the worktree root");
assertEquals("/wt/x", PeerLauncher.ideProjectPath("/wt/x", null),
"no module subdir ⇒ the worktree root itself");
assertEquals("/wt/x", PeerLauncher.ideProjectPath("/wt/x", " "),
"a blank module subdir ⇒ the worktree root itself");
}
@Test
void openInIdeIsNoOpAndNeverThrowsWhenCommandBlank() {
// A profile that opts into IDE MCP but sets no open command must not fail the spawn.
assertDoesNotThrow(() -> PeerLauncher.openInIde("/wt/x", null, LoggerFactory.getLogger("test")));
assertDoesNotThrow(() -> PeerLauncher.openInIde("/wt/x", " ", LoggerFactory.getLogger("test")));
}
@Test
void writeIdeOverlayDoesNothingWhenDotGitIsADirectory(@TempDir Path root) throws Exception {
Path worktree = Files.createDirectory(root.resolve("worktree"));
// The primary's real checkout has a `.git` DIRECTORY, not the worktree's `.git` FILE.
Files.createDirectories(worktree.resolve(".git"));
FakeHerdr herdr = new FakeHerdr();
launcher(herdr, ideProfile(null, "http://127.0.0.1:29170/index-mcp/streamable-http",
worktree.toString())).spawn();
assertFalse(Files.exists(worktree.resolve("CLAUDE.local.md")),
"the safety gate refuses to write into a non-worktree cwd (.git directory)");
}
@Test
void startRetriesWhileTheSeedShellBoots() {
// tab.create returns before the seed shell reaches its prompt; herdr refuses agent.start
@@ -136,7 +136,7 @@ class CompositePeerLauncherTest {
return new FleetConfig.Profile(profile, "http://gx00.gw:8000", "coder",
null, "BRIDGED_WORKER_TOKEN", List.of("claude"), "tab", "bridged-workers",
"w #{n}", null, null, null, null, null, null, null, null, null,
null, null, credentialId);
null, null, credentialId, null);
}
/**
@@ -494,4 +494,56 @@ class OpenCodeLauncherTest {
assertTrue(json.path("provider").isMissingNode(),
"without a baseUrl opencode resolves its own provider as before");
}
// --- CB-634: IDE Index MCP + guidance overlay (opencode does not read CLAUDE.local.md) -------
private static FleetConfig.Profile opencodeIdeCfg(String mcpUrl, String ideUrl, String cwd) {
return new FleetConfig.Profile("gemini", null, null, null, "BRIDGED_WORKER_TOKEN",
List.of("opencode"), "tab", "bridged-workers", "opencode: {model} #{n}", mcpUrl,
cwd, null, null, null, FleetConfig.Profile.KIND_OPENCODE,
null, null, null, null, null, null, ideUrl);
}
@Test
void ideMcpUrlAddsTheIntellijServerAndAnInstructionsRulesEntry(@TempDir Path root) throws Exception {
FakeHerdr herdr = new FakeHerdr();
Path cwd = Files.createDirectory(root.resolve("checkout"));
service(herdr, root, opencodeIdeCfg("http://127.0.0.1:8765/mcp",
"http://127.0.0.1:29170/index-mcp/streamable-http", cwd.toString())).spawn();
String cfgPath = startEnv(herdr).get("OPENCODE_CONFIG");
assertNotNull(cfgPath, "an IDE profile needs a config file");
JsonNode json = new ObjectMapper().readTree(Path.of(cfgPath).toFile());
JsonNode ide = json.path("mcp").path("intellij");
assertEquals("remote", ide.path("type").asText(),
"the IDE server uses the same remote shape as the bridge mount");
assertEquals("http://127.0.0.1:29170/index-mcp/streamable-http", ide.path("url").asText());
assertTrue(ide.path("enabled").asBoolean(), "the IDE server is enabled");
assertEquals("remote", json.path("mcp").path("fleet").path("type").asText(),
"the bridge mount still coexists with the IDE server");
// The instructions array gains an entry pointing at a real rules file pinning the worktree.
String rulesContent = null;
for (JsonNode n : json.path("instructions")) {
Path p = Path.of(n.asText());
if (p.getFileName().toString().equals("ide-rules.md")) {
rulesContent = Files.readString(p);
}
}
assertNotNull(rulesContent, "an ide-rules.md instructions entry is present");
assertTrue(rulesContent.contains("project_path: \"" + cwd + "\""),
"the rules pin every ide_* call to the worker's own cwd");
}
@Test
void noIdeServerWhenIdeMcpUrlUnset(@TempDir Path root) throws Exception {
FakeHerdr herdr = new FakeHerdr();
service(herdr, root, opencodeCfg("google/gemini-2.5-pro", "http://127.0.0.1:8765/mcp", null))
.spawn();
JsonNode json = new ObjectMapper()
.readTree(Path.of(startEnv(herdr).get("OPENCODE_CONFIG")).toFile());
assertTrue(json.path("mcp").path("intellij").isMissingNode(),
"no IDE server when ideMcpUrl is unset");
}
}
@@ -0,0 +1,173 @@
package dev.ltms.fleet.msg;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;
import org.testcontainers.containers.RabbitMQContainer;
import org.testcontainers.junit.jupiter.Testcontainers;
import org.testcontainers.utility.DockerImageName;
import java.util.List;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Contract test for {@link LeadMailbox} against a REAL broker — same approach as
* {@code AmqpReplyInboxContractTest}, which this mirrors: a Testcontainers RabbitMQ locally, or an
* externally-provisioned broker in CI via {@code AMQP_URI}. Tagged {@code contract} so it is
* excluded from {@code mvn test}/{@code mvn clean install} (which stay hermetic and need no
* Docker); run it with Docker present via {@code mvn test -Pcontract}.
*
* <p>Proves the mechanism this ticket adds: a {@link LeadMessage} published to
* {@code lead.<to>.inbox} is received with {@code from}/{@code to}/{@code content} intact, and
* {@link LeadMailbox#ack} removes it — the same publish→peek→ack roundtrip
* {@code AmqpReplyInboxContractTest} proves for {@link AmqpReplyInbox}, adapted to this class's
* single-owned-mailbox shape (no {@code own}/{@code release} — the mailbox for {@code selfCoordId}
* is owned the moment {@link LeadMailbox#open} returns).
*/
@Tag("contract")
// disabledWithoutDocker=false: on the CI path (AMQP_URI set) no container is started and the class
// must still run against the external broker even though the runner has no Docker.
@Testcontainers(disabledWithoutDocker = false)
class LeadMailboxTest {
private static final String EXTERNAL_URI = System.getenv("AMQP_URI");
private static final RabbitMQContainer BROKER =
new RabbitMQContainer(DockerImageName.parse("rabbitmq:3.13-management"));
private static final AtomicLong SEQ = new AtomicLong();
// No @Container: the JUnit 5 extension would force-start it even when AMQP_URI is set. Start it
// manually only on the local (no-external-broker) path; Ryuk reaps it on JVM exit.
@BeforeAll
static void startBrokerUnlessExternal() {
if (EXTERNAL_URI == null) {
BROKER.start();
}
}
private static String uri() {
if (EXTERNAL_URI != null) {
return EXTERNAL_URI;
}
// No trailing slash: an empty path is vhost "", which does not exist — omitting it selects
// the default vhost "/".
return "amqp://guest:guest@" + BROKER.getHost() + ":" + BROKER.getAmqpPort();
}
/** A fresh coord-id per test run so parallel/repeat runs never collide on the same queue. */
private static String coordId(String prefix) {
return prefix + "-" + System.nanoTime() + "-" + SEQ.incrementAndGet();
}
@Test
void publishThenPeekThenAckRoundTrip() throws Exception {
String to = coordId("lead-to");
String from = "lead-from";
try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) {
LeadMessage sent = new LeadMessage("m1", from, to, "hello peer lead");
inbox.publish(to, sent);
List<LeadMessage> got = awaitPeek(inbox);
assertEquals(1, got.size(), "the published message should be held for drain");
assertEquals("m1", got.getFirst().msgId());
assertEquals(from, got.getFirst().from());
assertEquals(to, got.getFirst().to());
assertEquals("hello peer lead", got.getFirst().content());
inbox.ack("m1");
assertTrue(inbox.peek().isEmpty(), "an acked message is dropped");
}
}
@Test
void duplicateMsgIdIsNotDoubleQueued() throws Exception {
String to = coordId("lead-dedup");
try (LeadMailbox inbox = LeadMailbox.open(uri(), to)) {
inbox.publish(to, new LeadMessage("dup", "lead-from", to, "first"));
awaitPeek(inbox);
inbox.publish(to, new LeadMessage("dup", "lead-from", to, "second")); // same msgId — no-op
Thread.sleep(500); // give any erroneous second delivery time to land
List<LeadMessage> got = inbox.peek();
assertEquals(1, got.size(), "a repeated msgId must not double-queue");
assertEquals("first", got.getFirst().content(), "the first payload wins");
}
}
@Test
void unackedMessageSurvivesRestartAndIsRedelivered() throws Exception {
String to = coordId("lead-durable");
// First "process life": publish, see it held, but crash before acking.
try (LeadMailbox first = LeadMailbox.open(uri(), to)) {
first.publish(to, new LeadMessage("persist-1", "lead-from", to, "survive me"));
assertEquals(1, awaitPeek(first).size());
// no ack — simulate a java -jar bounce with the message still pending
}
// Second "process life": a fresh connection owning the same mailbox must be redelivered it.
try (LeadMailbox second = LeadMailbox.open(uri(), to)) {
List<LeadMessage> got = awaitPeek(second);
assertEquals(1, got.size(), "an unacked persistent message is redelivered after restart");
assertEquals("persist-1", got.getFirst().msgId());
assertEquals("survive me", got.getFirst().content());
second.ack("persist-1");
}
// Third life: once acked, it is gone for good.
try (LeadMailbox third = LeadMailbox.open(uri(), to)) {
Thread.sleep(500);
assertTrue(third.peek().isEmpty(), "an acked message does not come back on the next restart");
}
}
@Test
void publishDoesNotRequireTheSenderToOwnTheTargetMailbox() throws Exception {
// CB-308 federation: a sender that never opened its own LeadMailbox for `to` can still
// publish to it — publish must not imply ownership. Only the owner ever consumes here.
String to = coordId("lead-federated");
String senderId = coordId("lead-sender");
try (LeadMailbox owner = LeadMailbox.open(uri(), to);
LeadMailbox sender = LeadMailbox.open(uri(), senderId)) {
sender.publish(to, new LeadMessage("m1", senderId, to, "from a federated peer"));
List<LeadMessage> got = awaitPeek(owner);
assertEquals(1, got.size(), "only the owner's mailbox should receive the message");
assertEquals(senderId, got.getFirst().from());
assertTrue(sender.peek().isEmpty(), "the sender must not also hold a copy — it never owns `to`");
}
}
@Test
void unroutablePublishReportsFailureNotSilentSuccess() throws Exception {
// Publish to a coord-id whose mailbox was never opened by anyone: the queue is never
// declared, so the default-exchange route to lead.<to>.inbox does not exist and the broker
// must return the publish.
String to = coordId("lead-nobody-home");
try (LeadMailbox sender = LeadMailbox.open(uri(), coordId("lead-sender"))) {
IllegalStateException ex = assertThrows(IllegalStateException.class,
() -> sender.publish(to, new LeadMessage("m1", "lead-from", to, "nobody home")));
assertTrue(ex.getMessage() != null && ex.getMessage().toLowerCase().contains("unroutable"),
"expected an unroutable-publish failure, got: " + ex.getMessage());
}
}
/** Poll peek until at least one message is held, or ~10s elapse (broker delivery is async). */
@SuppressWarnings("BusyWait")
private static List<LeadMessage> awaitPeek(LeadMailbox inbox) throws InterruptedException {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
List<LeadMessage> msgs = inbox.peek();
while (msgs.isEmpty() && System.nanoTime() < deadline) {
Thread.sleep(50);
msgs = inbox.peek();
}
return msgs;
}
}
@@ -252,7 +252,7 @@ class FleetAppTest {
FleetConfig.Profile wcfg = new FleetConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN", null,
"tab", "bridged-workers", "worker: {profile} #{n}", null,
null, null, null, null, null, null, null, 0, null, null, null);
null, null, null, null, null, null, null, 0, null, null, null, null);
Map<String, FleetConfig.Profile> profiles = Map.of(wcfg.profile(), wcfg);
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(
new AgentControl(herdr), new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
+16
View File
@@ -179,6 +179,22 @@ else
warn "This only matters once a profile points at llm.ltms.dev — harmless before that."
fi
# Third variable, same trap (CB-635). broker.uriEnv names the env var holding the AMQP URI, so the
# password stays out of bridged.yaml — but that moves the failure into the environment. If the
# variable is empty the daemon still starts: since #152 it warns and falls back to the in-memory
# reply inbox, so nothing crashes and replies simply stop surviving a restart. Only this check says
# so before the fact. Read the name out of the config so a renamed key cannot make the check lie.
BROKER_URI_ENV=$(sed -n 's/^[[:space:]]*uriEnv:[[:space:]]*\([A-Za-z_][A-Za-z0-9_]*\).*/\1/p' "$BRIDGED/bridged.yaml" | head -1)
if [ -z "$BROKER_URI_ENV" ]; then
ok "no broker.uriEnv configured — reply inbox is in-memory by design"
elif zsh -lc "[ -n \"\${$BROKER_URI_ENV:-}\" ]" 2>/dev/null; then
ok "$BROKER_URI_ENV (broker.uriEnv) resolves in a login shell"
else
warn "$BROKER_URI_ENV (broker.uriEnv) is EMPTY in a login shell."
warn "The daemon will start and fall back to the IN-MEMORY reply inbox."
warn "Replies stop surviving a restart — a held report is lost, not delayed."
fi
if [ "$CHECK_ONLY" = 1 ]; then
say "--check: nothing changed"
exit 0