diff --git a/CLAUDE.md b/CLAUDE.md index 5dcf182..a404263 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -112,7 +112,8 @@ the merge — and merging on a reviewer's word is delegating it by proxy. | Delegate (blocking) | `fleet_send{sessionId, content}` | | Delegate (long task) | `fleet_send{sessionId, content, wait:false}` → ticket → `fleet_poll{ticket}` | | Answer a member's `fleet_ask` | `fleet_send{turnId, content}` — **not** `sessionId` | -| Message a **peer lead** | `fleet_send{sessionId: , content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task | +| Message a **peer lead** on this host | `fleet_send{sessionId: , content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task | +| Message a **peer lead** on another daemon or host | `fleet_send{coordId: , content}` — needs a `coordinator:` block; your own coord-id is in `fleet_list`. Coordination only, **never** a task | | Answer a peer lead that messaged you | `fleet_reply{content}` — the one case a lead replies | | Collect a held reply | `fleet_poll{target}` · then `fleet_ack{target, msgId}` | | Tear down a member | `fleet_stop{paneId}` | diff --git a/bridged/fleetd.example.yaml b/bridged/fleetd.example.yaml index eea90c6..8c47d03 100644 --- a/bridged/fleetd.example.yaml +++ b/bridged/fleetd.example.yaml @@ -126,6 +126,32 @@ 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. +# autoCompactWindow → opt-in, default off. A bounded token window that forces a spawned member to +# compact its context instead of running on the backend's own default and dying +# mid-turn (losing its fleet_reply — the whole point of the turn — with it). +# Validated at config load to [100000, 1000000] — the band Claude Code's own +# --autocompact flag accepts. +# CROSS-BACKEND SEMANTICS DIFFER: on claude-code this is a launch-time +# `--autocompact ` flag — the member compacts AT this window. opencode +# has no equivalent flag (it only forces `compaction.auto: true`, unconditionally, +# already), so this is instead applied as the model's `limit.context` in the +# generated opencode.json — the member compacts WITHIN this window, not exactly +# at it — and only when this profile's `model:` is in `provider/model` form; if it +# isn't, bridged logs a WARN naming the profile rather than silently doing nothing. # 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 +161,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 +264,11 @@ 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 + # autoCompactWindow: 250000 # opt-in: bound member context; claude-code compacts AT this, opencode within it (model limit.context) 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 +667,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..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 diff --git a/bridged/src/main/java/dev/ltms/fleet/Fleetd.java b/bridged/src/main/java/dev/ltms/fleet/Fleetd.java index 65dd744..88fc751 100644 --- a/bridged/src/main/java/dev/ltms/fleet/Fleetd.java +++ b/bridged/src/main/java/dev/ltms/fleet/Fleetd.java @@ -31,6 +31,9 @@ import dev.ltms.fleet.mcp.LsofPeerPidLookup; import dev.ltms.fleet.mcp.LsofProcessCwdLookup; import dev.ltms.fleet.msg.AmqpReplyInbox; import dev.ltms.fleet.msg.InMemoryReplyInbox; +import dev.ltms.fleet.msg.LeadChannel; +import dev.ltms.fleet.msg.LeadCoordLoop; +import dev.ltms.fleet.msg.LeadMailbox; import dev.ltms.fleet.msg.MessageService; import dev.ltms.fleet.msg.Rendezvous; import dev.ltms.fleet.msg.ReplyInbox; @@ -60,6 +63,7 @@ import java.util.Map; import java.util.Objects; import java.util.Set; import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Function; @@ -79,6 +83,14 @@ public final class Fleetd { /** CB-504: how long to wait at startup for herdr's socket before serving degraded. */ private static final long HERDR_WAIT_SECONDS = 30; + /** + * CB-637: how often the lead coordination loop looks for peer messages. A few seconds — slow + * enough that an idle fleet is not polling a broker in a tight loop, fast enough that a peer + * lead's message is not left sitting once the local lead reaches a turn boundary. The mailbox + * pushes into the loop's held set on its own consumer thread, so this interval bounds only the + * pane delivery, never the receive. + */ + private static final long LEAD_COORD_INTERVAL_MS = 3_000L; private static final long HERDR_WAIT_POLL_MILLIS = 500; /** @@ -375,6 +387,12 @@ public final class Fleetd { // unusable), bridged stays soft-state on the in-memory inbox. The AMQP inbox owns a broker // connection, so keep the reference to close it in the ordered shutdown hook. final ReplyInbox replyInbox = selectReplyInbox(cfg.broker(), System.getenv(), AmqpReplyInbox::open); + // CB-637: this daemon's lead-to-lead mailbox on the SHARED coordination vhost — a separate + // broker from the reply inbox by design (see FleetConfig.Coordinator). Absent a coordinator: + // block this is null and every lead path below is simply not wired, which is exactly the + // behaviour before this ticket. It owns a broker connection, so keep the reference for the + // ordered shutdown hook. + final LeadMailbox leadMailbox = openLeadMailbox(cfg.coordinator(), System.getenv(), LeadMailbox::open); // CB-307: learn the primary's terminal from orchestration tool calls (or pin from config). // The pin also feeds CallerResolver below: a primary running inside a herdr pane would // otherwise resolve as a worker and be refused every orchestration tool. @@ -503,7 +521,27 @@ public final class Fleetd { new FleetMcp.QuarantineSource(profile -> { var configured = config.get().profiles().get(profile); return configured == null ? null : configured.effectiveCredentialId(); - }, quarantine)); + }, quarantine), + leadMailbox); + + // CB-637: the receive half. Only constructed when a lead mailbox actually opened — with no + // coordinator (or an unreachable one) there is nothing to deliver, so no scheduler is + // created and no thread runs. It reads the SAME live lead supplier the injector's + // deliverability gate does, so a lead found by the tab scan after startup is reachable + // without a restart. + final LeadCoordLoop leadCoordLoop; + final ScheduledExecutorService leadCoordSchedulerRef; + if (leadMailbox != null) { + var leadCoordScheduler = Executors.newSingleThreadScheduledExecutor(r -> + Thread.ofVirtual().name("bridge-leadcoord-").unstarted(r)); + leadCoordLoop = new LeadCoordLoop(leadMailbox, agents, leads, leadCoordScheduler, + LEAD_COORD_INTERVAL_MS); + leadCoordLoop.start(); + leadCoordSchedulerRef = leadCoordScheduler; + } else { + leadCoordLoop = null; + leadCoordSchedulerRef = null; + } // 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. @@ -524,6 +562,8 @@ public final class Fleetd { messages.close(); pushLoop.close(); if (heartbeat != null) heartbeat.close(); // CB-551: stop the idle-lead heartbeat scheduler + if (leadCoordLoop != null) leadCoordLoop.close(); // CB-637: stop delivering peer-lead messages + if (leadCoordSchedulerRef != null) leadCoordSchedulerRef.shutdownNow(); if (healthMonitor != null) healthMonitor.stop(); if (configWatcher != null) configWatcher.stop(); // CB-559: stop polling the config file mcp.close(); @@ -536,6 +576,15 @@ public final class Fleetd { log.debug("reply inbox close: {}", e.toString()); } } + // CB-637: the coordination connection goes with it — after the loop that reads it has + // stopped, so no tick can be mid-ack against a closed channel. + if (leadMailbox != null) { + try { + leadMailbox.close(); + } catch (Exception e) { + log.debug("lead mailbox close: {}", e.toString()); + } + } herdr.close(); })); @@ -575,6 +624,66 @@ public final class Fleetd { ReplyInbox open(String uri, int prefetch); } + /** Injection seam for {@link #openLeadMailbox}: production binds {@link LeadMailbox#open}. */ + @FunctionalInterface + interface LeadMailboxOpener { + LeadMailbox open(String uri, String selfCoordId, int prefetch); + } + + /** + * CB-637: open this daemon's lead-to-lead mailbox, or return {@code null} to leave the feature + * off. Package-private and env-injected for the same reason as {@link #selectReplyInbox}: the + * selection is then testable without a broker or a mutable process environment. + * + *

Every "off" path returns {@code null}, and each says why at the level it deserves: + * + *

    + *
  • no {@code coordinator:} block — silent. Lead coordination is opt-in; an operator who + * never configured it does not need to be told it is off on every boot.
  • + *
  • a block whose {@code uriEnv} does not resolve — INFO, the same "you moved to the secret + * store and the variable is not there" case {@code selectReplyInbox} warns about.
  • + *
  • a configured broker but no {@code selfId} — WARN. This one is a half-finished config: a + * mailbox is named after the coord-id that owns it, so with no id there is no queue to own + * and no {@code from} to send as. Loud, because the operator plainly intended the feature.
  • + *
  • the broker refuses at boot — WARN, and carry on. Mirrors {@code openAmqpOrFallback}: a + * coordination broker that is down must never take a whole fleet's daemon with it, and the + * fleet still works exactly as it did before this feature existed.
  • + *
+ */ + static LeadMailbox openLeadMailbox(FleetConfig.Coordinator coordinator, Map env, + LeadMailboxOpener opener) { + if (coordinator == null) { + return null; // opt-in: nothing configured, nothing to say + } + String uri = coordinator.effectiveUri(env); + if (uri == null) { + log.info("lead coordination: OFF — coordinator{} has no usable broker uri", + coordinator.uriEnv() == null ? "" : ".uriEnv=" + coordinator.uriEnv()); + return null; + } + if (coordinator.selfId() == null || coordinator.selfId().isBlank()) { + log.warn("coordinator.selfId is unset — lead coordination is OFF. A lead mailbox is the " + + "queue named after the coord-id that owns it, so with no id there is nothing to " + + "own and no sender identity to publish as. Set coordinator.selfId to a name that " + + "is unique across every daemon sharing {} and restart bridged.", + stripCredentials(uri)); + return null; + } + try { + LeadMailbox mailbox = opener.open(uri, coordinator.selfId(), coordinator.prefetchOrDefault()); + log.info("lead coordination: ON as coord-id {} (prefetch={})", + coordinator.selfId(), coordinator.prefetchOrDefault()); + return mailbox; + } catch (IllegalStateException e) { + log.warn("cannot reach the AMQP coordination broker ({}) — lead-to-lead messaging is OFF " + + "for this process lifetime. fleet_send{{coordId}} will report it as " + + "not configured, and peer messages already queued stay on the broker until " + + "a restart picks them up. Reason: {}", + stripCredentials(uri), reasonOf(e)); + return null; + } + } + /** * CB-151/152: pick the reply inbox. A usable broker — a literal {@code uri}, or a {@code * uriEnv} whose variable resolves (both read from {@code env}) — selects the durable AMQP inbox. diff --git a/bridged/src/main/java/dev/ltms/fleet/config/ConfigRef.java b/bridged/src/main/java/dev/ltms/fleet/config/ConfigRef.java index 651c6a8..f230e3b 100644 --- a/bridged/src/main/java/dev/ltms/fleet/config/ConfigRef.java +++ b/bridged/src/main/java/dev/ltms/fleet/config/ConfigRef.java @@ -276,6 +276,9 @@ public final class ConfigRef implements Supplier { && 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()) diff --git a/bridged/src/main/java/dev/ltms/fleet/config/FleetConfig.java b/bridged/src/main/java/dev/ltms/fleet/config/FleetConfig.java index ea84552..5a34371 100644 --- a/bridged/src/main/java/dev/ltms/fleet/config/FleetConfig.java +++ b/bridged/src/main/java/dev/ltms/fleet/config/FleetConfig.java @@ -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 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 profiles, Guard guard, @@ -275,6 +294,19 @@ public record FleetConfig( * profile that does not opt in. Read live off the current config, so it is * HOT: a change takes effect on the next exhaustion classification / spawn, * no restart needed. + * @param autoCompactWindow opt-in per-profile token window that forces a spawned member to + * auto-compact its context at (Claude Code) or within (opencode) a bound the + * operator chooses, instead of the backend's own default. {@code null} (the + * default) leaves today's behaviour exactly — opencode already forces + * {@code compaction.auto: true} unconditionally (CB-523) but has no absolute + * window, and Claude Code has neither. When set, validated at config load to + * {@code [100000, 1000000]} — the band Claude Code's own {@code --autocompact + * } flag accepts. The two backends honour it differently: Claude Code + * compacts AT this window (a launch-time {@code --autocompact} flag); + * opencode has no such knob, so this is applied as the model's + * {@code limit.context} instead, which bounds the window opencode compacts + * within, and only when {@code model} resolves to a + * {@code provider/model} pair. */ @JsonIgnoreProperties(ignoreUnknown = true) public record Profile(String profile, String baseUrl, String model, @@ -289,7 +321,11 @@ public record FleetConfig( Integer maxLoad, Boolean subscription, String exhaustedPattern, - String credentialId) { + String credentialId, + String ideMcpUrl, + String ideProjectDir, + String ideOpenCommand, + Integer autoCompactWindow) { /** Peer kind spawned by {@link dev.ltms.fleet.member.ClaudeCodeLauncher} (the default). */ public static final String KIND_CLAUDE_CODE = "claude-code"; @@ -348,6 +384,21 @@ 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; + // Opt-in per profile, default off (null). No clamping here — unlike ideMcpUrl/ideProjectDir + // there is no blank-string form to normalize (it's an Integer), and the [100000, 1000000] + // range is enforced eagerly at config load (rejectAutoCompactWindowOutOfRange), naming the + // profile, rather than silently clamped here. A profile that never sets it keeps null. } /** @@ -392,7 +443,47 @@ 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, autoCompactWindow); + } + + /** + * Backward-compatible constructor without the {@code autoCompactWindow} field — the profile + * leaves auto-compaction at the backend's own default (opencode's unconditional + * {@code compaction.auto: true}, or Claude Code's built-in threshold), exactly as before this + * key existed. This is the shape the canonical constructor had before the field was added — + * every pre-existing Java call site (and any YAML that omits the key) keeps compiling and + * behaving identically; Jackson binds the canonical (longest) constructor, so YAML omitting + * {@code autoCompactWindow:} still lands here as {@code null} via that path, not this one. + */ + public Profile(String profile, String baseUrl, String model, + String configDir, String tokenEnv, List argv, + String placement, String workspace, String tabLabel, String mcpUrl, + String cwd, List parityOverlay, String gitTokenEnv, String gitHostEnv, + String kind, Map env, Float weight, Integer maxLoad, + Boolean subscription, String exhaustedPattern, String credentialId, + String ideMcpUrl, String ideProjectDir, String ideOpenCommand) { + this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel, + mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad, + subscription, exhaustedPattern, credentialId, ideMcpUrl, ideProjectDir, ideOpenCommand, + null); + } + + /** + * 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 argv, + String placement, String workspace, String tabLabel, String mcpUrl, + String cwd, List parityOverlay, String gitTokenEnv, String gitHostEnv, + String kind, Map 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 +532,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 +580,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 +709,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. + * + *

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..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 not 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 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 +1319,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) { @@ -1150,6 +1330,7 @@ public record FleetConfig( warnUnknownTopLevelKeys(yaml, path); rejectDuplicateMemberSlots(yaml); rejectNegativeMaxLoad(yaml); + rejectAutoCompactWindowOutOfRange(yaml); rejectUnknownKind(yaml); rejectUnknownAuthMode(yaml); rejectUnknownPlacement(yaml); @@ -1452,6 +1633,52 @@ public record FleetConfig( } } + /** Lowest {@code autoCompactWindow} Claude Code's {@code --autocompact } flag accepts. */ + static final int AUTO_COMPACT_WINDOW_MIN = 100_000; + /** Highest {@code autoCompactWindow} Claude Code's {@code --autocompact } flag accepts. */ + static final int AUTO_COMPACT_WINDOW_MAX = 1_000_000; + + /** + * Reject a profile whose {@code autoCompactWindow:} is set but outside the token band Claude + * Code's own {@code --autocompact } flag accepts (100k–1M), naming both the profile and + * the value. + * + *

Unset/{@code null} means "off" and passes silently — today's behaviour for every profile + * that does not opt in (see {@link Profile#autoCompactWindow()}). A profile that DOES set the key + * is validated eagerly, at config load, rather than failing later when Claude Code itself refuses + * the launch flag on spawn — the same "fail loud at load, not lazily at first spawn" reasoning as + * {@link #rejectNegativeMaxLoad} and {@link #rejectUnknownPlacementPolicy}. + * + * @param yaml the raw config text + * @throws IllegalStateException when any profile's {@code autoCompactWindow} is set and outside + * {@code [100000, 1000000]} + */ + static void rejectAutoCompactWindowOutOfRange(String yaml) { + Map raw; + try { + raw = YAML.readValue(yaml, Map.class); + } catch (IOException | IllegalArgumentException e) { + return; // a malformed file is reported by the real parse, not here + } + if (raw == null || !(raw.get("profiles") instanceof Map profiles)) { + return; + } + List bad = profiles.entrySet().stream() + .filter(e -> e.getValue() instanceof Map p + && p.get("autoCompactWindow") instanceof Number n + && (n.doubleValue() < AUTO_COMPACT_WINDOW_MIN || n.doubleValue() > AUTO_COMPACT_WINDOW_MAX)) + .map(e -> String.valueOf(e.getKey())) + .sorted() + .toList(); + if (!bad.isEmpty()) { + throw new IllegalStateException("refusing to start: profile(s) [" + String.join(", ", bad) + + "] set autoCompactWindow outside [" + AUTO_COMPACT_WINDOW_MIN + ", " + + AUTO_COMPACT_WINDOW_MAX + "] — Claude Code's --autocompact flag accepts only " + + "that band of tokens; omit the key to leave auto-compaction at each backend's " + + "own default."); + } + } + /** The peer kinds this build has an adapter for — {@link Profile#kind()}'s only valid values. */ private static final Set KNOWN_KINDS = Set.of(Profile.KIND_CLAUDE_CODE, Profile.KIND_OPENCODE); @@ -1708,9 +1935,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); } /** diff --git a/bridged/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java b/bridged/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java index fe67ef2..7e5ef03 100644 --- a/bridged/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java +++ b/bridged/src/main/java/dev/ltms/fleet/mcp/FleetMcp.java @@ -11,6 +11,8 @@ import dev.ltms.fleet.metrics.FleetMetrics; import dev.ltms.fleet.metrics.Metrics; import dev.ltms.fleet.inject.MemberPresence; import dev.ltms.fleet.herdr.HerdrException; +import dev.ltms.fleet.msg.LeadChannel; +import dev.ltms.fleet.msg.LeadMessage; import dev.ltms.fleet.msg.MessageService; import dev.ltms.fleet.msg.Rendezvous; import dev.ltms.fleet.peer.PeerUnreachableException; @@ -37,6 +39,7 @@ import java.util.List; import java.util.Map; import java.util.Objects; import java.util.Set; +import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.function.BiFunction; import java.util.function.Function; @@ -97,6 +100,8 @@ public final class FleetMcp { private final CapacitySource capacity; private final HealthCoverageSource healthCoverage; private final QuarantineSource quarantine; + /** CB-637: this daemon's lead-to-lead channel; {@code null} when no coordinator is configured. */ + private final LeadChannel leadChannel; /** Capacity facts used by {@code fleet_list}; production must supply the placement live count. */ public record CapacitySource(Function liveCount, Function maxLoad, @@ -131,6 +136,21 @@ public final class FleetMcp { ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry, CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine) { + this(messages, workers, sessions, identity, presence, primaryRegistry, callers, metrics, capacity, + healthCoverage, quarantine, null); + } + + /** + * As above, with this daemon's lead-to-lead channel (CB-637). {@code leadChannel} is + * {@code null} whenever no {@code coordinator:} block is configured or its broker could not be + * reached at boot — cross-daemon lead messaging is simply off, and {@code fleet_send{coordId}} + * says so rather than failing obscurely. + */ + public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions, + ConnectionIdentity identity, MemberPresence presence, PrimaryRegistry primaryRegistry, + CallerResolver callers, Metrics metrics, CapacitySource capacity, HealthCoverageSource healthCoverage, + QuarantineSource quarantine, LeadChannel leadChannel) { + this.leadChannel = leadChannel; this.capacity = capacity; this.quarantine = Objects.requireNonNull(quarantine, "quarantine"); this.healthCoverage = healthCoverage; @@ -177,6 +197,13 @@ public final class FleetMcp { String target = str(a, "sessionId"); String content = str(a, "content"); String turnId = str(a, "turnId"); + String coordId = str(a, "coordId"); + if (coordId != null && !coordId.isBlank()) { + // CB-637: a peer LEAD on another daemon, addressed by coord-id over the shared + // coordination broker. Checked before the turnId branch so a call that sets both + // is rejected as the conflict it is, rather than silently taking one route. + return sendToLead(leadChannel, coordId, content, target, turnId); + } if (turnId != null && !turnId.isBlank()) { // Answering a worker's fleet_ask (CB-205): resolve its blocked question and // block for the worker's reply as it resumes the same turn. This is the same @@ -257,7 +284,8 @@ public final class FleetMcp { if (denied != null) return denied; return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, callers == null ? Map.of() : callers.leads(), - callerTerminal(exchange)); + callerTerminal(exchange), + leadChannel == null ? null : leadChannel.selfCoordId()); }; BiFunction stopHandler = (exchange, req) -> { @@ -585,6 +613,54 @@ public final class FleetMcp { return null; } + /** + * {@code fleet_send} carrying a {@code coordId} (CB-637): a message to a PEER LEAD, published to + * that lead's durable mailbox on the shared coordination broker. This is the only lead→lead path + * that crosses hosts — the existing pane-injection route can only reach a lead whose herdr socket + * this daemon shares. + * + *

{@code coordId} is mutually exclusive with {@code sessionId} and {@code turnId}: those two + * address a worker session owned by this daemon, a coord-id addresses a lead owned by + * another one, and there is no sensible reading of a call that sets both. Rejected by name rather + * than resolved by precedence, so a caller that meant the other route learns it instead of having + * its message quietly go somewhere else. + * + *

The publish is synchronous and confirmed by the broker, so the result is a real delivery + * receipt rather than a hopeful one. Its failure — nobody owns {@code coordId}'s mailbox, the + * broker nacked, or the confirm timed out — arrives as {@link IllegalStateException} and is + * turned into a tool error naming the coord-id. It is never allowed to escape as a crash: an + * unreachable peer is an ordinary outcome of addressing a fleet you do not control. + * + * @param leadChannel this daemon's channel, or {@code null} when no coordinator is configured + */ + static McpSchema.CallToolResult sendToLead(LeadChannel leadChannel, String coordId, String content, + String sessionId, String turnId) { + if (!isBlank(sessionId) || !isBlank(turnId)) { + String conflict = !isBlank(sessionId) ? "sessionId" : "turnId"; + return error("coordId and " + conflict + " are mutually exclusive: coordId addresses a peer " + + "LEAD on another daemon over the coordination broker, while " + conflict + + " addresses a worker session on this one. Pass exactly one."); + } + if (isBlank(content)) { + return error("content is required"); + } + if (leadChannel == null) { + return error("lead coordination is not configured (no coordinator: block) — cannot send to " + + "peer lead \"" + coordId + "\". Add a coordinator: block with a shared broker uri " + + "and this daemon's selfId, then restart bridged."); + } + LeadMessage msg = new LeadMessage(UUID.randomUUID().toString(), leadChannel.selfCoordId(), + coordId, content); + try { + leadChannel.publish(coordId, msg); + } catch (IllegalStateException e) { + return error("cannot deliver to peer lead \"" + coordId + "\": " + e.getMessage() + + ". Check that a daemon is running with coordinator.selfId=\"" + coordId + + "\" and is connected to the same coordination broker."); + } + return text("delivered to peer lead " + coordId + " (msgId " + msg.msgId() + ")"); + } + /** {@code fleet_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)) { @@ -870,6 +946,22 @@ public final class FleetMcp { static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, CapacitySource capacity, HealthCoverageSource healthCoverage, QuarantineSource quarantine, Map leads, String selfTerm) { + return listFleet(workers, sessions, messages, capacity, healthCoverage, quarantine, leads, selfTerm, null); + } + + /** + * As above, additionally reporting this daemon's own lead coordination id (CB-637) when one is + * configured and its channel opened. There is no peer-discovery surface yet — a lead addresses a + * peer by a coord-id it was told — so this row exists to answer the one question the operator + * cannot answer any other way: what is MY coord-id, the one a peer must use to reach me. It is + * omitted entirely when no coordinator is configured, so an ordinary fleet's output is unchanged. + * + * @param selfCoordId this daemon's coord-id, or {@code null} when lead coordination is off + */ + static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages, + CapacitySource capacity, HealthCoverageSource healthCoverage, + QuarantineSource quarantine, Map leads, String selfTerm, + String selfCoordId) { try { Map live = workers.list().stream() .map(Agent.class::cast) @@ -888,6 +980,9 @@ public final class FleetMcp { Map result = new LinkedHashMap<>(); result.put("leads", leadRows); result.put("members", out); result.put("healthCoverage", healthCoverage.value().get()); + if (selfCoordId != null && !selfCoordId.isBlank()) { + result.put("coordinator", Map.of("selfId", selfCoordId, "configured", true)); + } if (capacity.available()) result.put("capacity", profiles.stream() .map(profile -> capacityView(profile, capacity.liveCount(), capacity.maxLoad(), roster, messages, capacity.clock().getAsLong(), quarantine)).toList()); @@ -1015,7 +1110,9 @@ public final class FleetMcp { "Delegate a task to a worker session. By default blocks until the worker replies and " + "returns its reply (or a 'still working / queued' note on timeout). Pass wait:false " + "for a long task to return a ticket immediately, then poll it with fleet_poll. To " - + "answer a worker's fleet_ask, pass its turnId (with content) instead of sessionId.", + + "answer a worker's fleet_ask, pass its turnId (with content) instead of sessionId. " + + "To message a PEER LEAD on another daemon — possibly another host — pass its " + + "coordId instead; that is coordination, never a task.", objectSchema(Map.of( "sessionId", stringProp("The worker session id (herdr terminal_id) to delegate to"), "content", stringProp("The task/message to send to the worker (or your answer, with turnId)"), @@ -1023,7 +1120,11 @@ public final class FleetMcp { "wait", Map.of("type", "boolean", "description", "Block for the reply (default true); false returns a ticket to poll"), "turnId", stringProp("When answering a worker's fleet_ask, its question turnId — " - + "routes your answer back into the same turn (omit for a normal delegation)")), + + "routes your answer back into the same turn (omit for a normal delegation)"), + "coordId", stringProp("A peer LEAD's coordination id — delivers content to that " + + "lead's durable mailbox on the shared coordination broker, which works " + + "across hosts. Mutually exclusive with sessionId and turnId. Your own " + + "coordId is reported by fleet_list.")), List.of("content"))); } diff --git a/bridged/src/main/java/dev/ltms/fleet/member/ClaudeCodeLauncher.java b/bridged/src/main/java/dev/ltms/fleet/member/ClaudeCodeLauncher.java index c972731..1a58907 100644 --- a/bridged/src/main/java/dev/ltms/fleet/member/ClaudeCodeLauncher.java +++ b/bridged/src/main/java/dev/ltms/fleet/member/ClaudeCodeLauncher.java @@ -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 workerEnv = baseEnv(cfg); if (onSubscription) { // CB-542 belt-and-braces: on the subscription path no guard vets these two keys, and the @@ -237,7 +248,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher { // has neither MCP nor a charter — session flags must be added into a list we own. List argv = mutableArgv(argvWithFleet(cfg, spec)); String agentSessionId = applySessionIdentity(argv, spec.sessionName(), spec.resumeSessionId()); - return new Launch(workerEnv, argvWithModel(argv, cfg), agentSessionId); + return new Launch(workerEnv, argvWithAutoCompact(argvWithModel(argv, cfg), cfg), agentSessionId); } /** @@ -292,26 +303,31 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher { */ private List 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 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 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). + * + *

Safety gate: the overlay is written ONLY when {@code cwd/.git} is a + * regular file — a provisioned worktree keeps a {@code .git} FILE holding a + * {@code gitdir: } 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. + * + *

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 + // /worktrees/, 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; @@ -370,6 +463,30 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher { return withModel; } + /** + * Pin a bounded auto-compaction window on the command line via {@code --autocompact }, + * opt-in per profile (CB-634's sibling ticket: a member that runs out of context dies mid-turn + * and its {@code fleet_reply} — the whole point of the turn — is lost with it; opencode already + * forces {@code compaction.auto: true} unconditionally, CB-523, but Claude Code has no equivalent + * and runs at the backend's own default window). + * + *

Mirrors {@link #argvWithModel}: appended after it, so it survives the {@code ccs } + * wrapper the same way {@code --model} does, and outranks env/settings and the operator's own + * {@code argv}. Verified: {@code claude 2.1.241 --help} lists {@code --autocompact } + * (either the literal {@code auto}, or an integer 100k–1M) — {@link FleetConfig#load} rejects a + * configured value outside that band before this ever runs, so the flag Claude Code receives here + * is always in range. + */ + private static List argvWithAutoCompact(List argv, FleetConfig.Profile cfg) { + if (cfg.autoCompactWindow() == null) { + return argv; + } + List withAutoCompact = mutableArgv(argv); + withAutoCompact.add("--autocompact"); + withAutoCompact.add(String.valueOf(cfg.autoCompactWindow())); + return withAutoCompact; + } + // --- Agent-returning convenience spawns (used by callers/tests that want the herdr Agent) --- /** Spawn a worker for the default profile in the resolved default cwd. */ diff --git a/bridged/src/main/java/dev/ltms/fleet/member/OpenCodeLauncher.java b/bridged/src/main/java/dev/ltms/fleet/member/OpenCodeLauncher.java index ff9b31b..673284a 100644 --- a/bridged/src/main/java/dev/ltms/fleet/member/OpenCodeLauncher.java +++ b/bridged/src/main/java/dev/ltms/fleet/member/OpenCodeLauncher.java @@ -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 opencode — 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,20 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher { @Override protected Launch buildLaunch(FleetConfig.Profile cfg, LaunchSpec spec) { Map 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()); + // autoCompactWindow's opencode lever (limit.context) only targets a specific provider/model + // entry, so it needs model: in "provider/model" form. A profile that opts in without that + // shape gets no silent no-op — log it, once, here, whether or not writeConfig ends up running. + boolean wantsContextLimit = cfg.autoCompactWindow() != null && splitProviderModel(cfg.model()) != null; + if (cfg.autoCompactWindow() != null && !wantsContextLimit) { + log.warn("profile '{}' sets autoCompactWindow but model '{}' is not \"/\" " + + "form — opencode's per-model context limit could not be applied for this profile", + cfg.profile(), cfg.model()); + } + // A config file is needed for the bridge MCP mount, a member charter, the IDE MCP (+ its + // guidance overlay, CB-634), a pinned endpoint (CB-508), or a resolvable autoCompactWindow. + if (cfg.hasMcp() || cfg.hasIdeMcp() || spec.charter() != null || hasCustomProvider(cfg) + || wantsContextLimit) { + workerEnv.put("OPENCODE_CONFIG", writeConfig(cfg, spec.charter(), spec.cwd()).toString()); } applyGitToken(workerEnv, cfg); List argv = argvWithResume(argvWithModel(argvWithAuto(cfg), cfg), spec.resumeSessionId()); @@ -290,7 +306,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,15 +337,42 @@ 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); } + if (cfg.autoCompactWindow() != null) { + applyContextLimit(root, cfg); + } Path cfgFile = dir.resolve("opencode.json"); // Built with Jackson rather than string concatenation: the provider block is nested and @@ -374,18 +417,68 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher { * default gateway — a worker quietly talking to the wrong endpoint is the failure this avoids. */ private static String[] splitModelSelector(FleetConfig.Profile cfg) { - String model = cfg.model(); - int slash = model == null ? -1 : model.indexOf('/'); - if (model == null || model.isBlank() || slash <= 0 || slash == model.length() - 1) { + String[] parts = splitProviderModel(cfg.model()); + if (parts == null) { throw new IllegalArgumentException( "profile " + cfg.profile() + " sets baseUrl (a pinned opencode endpoint) so" + " model: must be \"/\", e.g." + " \"local-vllm/deepseek-v4-flash\"; got " - + (model == null ? "null" : '"' + model + '"')); + + (cfg.model() == null ? "null" : '"' + cfg.model() + '"')); + } + return parts; + } + + /** + * Split {@code model} into its {@code provider} and {@code model} halves, or {@code null} when + * it is not in that shape (unset/blank, or no non-trailing {@code /}). Unlike + * {@link #splitModelSelector}, non-throwing — callers that only *optionally* need the split + * (autoCompactWindow's context-limit application) use this to fall back to a WARN rather than an + * exception, since a profile without {@code baseUrl} is not required to name a provider/model. + */ + private static String[] splitProviderModel(String model) { + int slash = model == null ? -1 : model.indexOf('/'); + if (model == null || model.isBlank() || slash <= 0 || slash == model.length() - 1) { + return null; } return new String[]{model.substring(0, slash), model.substring(slash + 1)}; } + /** + * Apply the per-profile {@code autoCompactWindow} as opencode's per-model context limit. + * + *

opencode has no absolute "compact at N tokens" knob — its {@code compaction} block only + * exposes {@code auto}/{@code prune}/{@code reserved}/{@code tail_turns}/ + * {@code preserve_recent_tokens} — so the real lever is the model's own + * {@code provider.

.models..limit.context}, which bounds the window opencode compacts + * within rather than compacting exactly AT it the way Claude Code's {@code --autocompact} + * does. + * + *

Uses get-or-create nodes ({@code withObject}) at every level so this MERGES with any provider + * block {@link #addCustomProvider} already wrote for a custom-provider (pinned-endpoint) profile — + * it must never overwrite that block's {@code npm}/{@code name}/{@code options}. For a gateway + * profile (no {@code baseUrl}, so no prior provider block) this writes a partial + * {@code provider.

.models..limit} override, which opencode merges over its own built-in + * provider definition. + * + *

opencode's {@code limit} schema requires both {@code context} and {@code output}; there is no + * independent signal for the latter here, so 16384 is written as a safe default (documented in + * {@code fleetd.example.yaml}). + * + *

Silently does nothing when {@code model:} is not in {@code provider/model} form — a warning + * for that case is already logged once in {@code buildLaunch}, so this stays quiet rather than + * duplicating it. + */ + private static void applyContextLimit(ObjectNode root, FleetConfig.Profile cfg) { + String[] parts = splitProviderModel(cfg.model()); + if (parts == null) { + return; + } + ObjectNode limit = root.withObject("provider").withObject(parts[0]) + .withObject("models").withObject(parts[1]).withObject("limit"); + limit.put("context", cfg.autoCompactWindow()); + limit.put("output", 16384); + } + /** * The OpenAI-compatible base URL for {@code baseUrl}. A bare {@code host:port} gets {@code /v1} * appended (where these servers put the API); a URL that already carries a path is taken as-is, diff --git a/bridged/src/main/java/dev/ltms/fleet/msg/LeadChannel.java b/bridged/src/main/java/dev/ltms/fleet/msg/LeadChannel.java new file mode 100644 index 0000000..07664a9 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/fleet/msg/LeadChannel.java @@ -0,0 +1,40 @@ +package dev.ltms.fleet.msg; + +import java.util.List; + +/** + * The lead-to-lead message channel this daemon speaks, as its callers need it — one lead's own + * mailbox: publish to a peer's coord-id, look at what has arrived for me, and ack what I have + * delivered. + * + *

Extracted from {@link LeadMailbox} purely as a seam. {@code LeadMailbox} is the one production + * implementation and owns a live AMQP connection, so a test that wanted to exercise the routing in + * {@code FleetMcp} or the delivery in {@link LeadCoordLoop} would have had to stand up a broker — + * which is exactly the kind of test that gets tagged {@code contract} and then does not run. With + * this interface both of those are hermetic: they inject a fake channel and assert on what was + * published, peeked and acked. + * + *

Note what is not here: {@code drain()} and {@code close()}. Draining is a convenience + * over peek+ack that no caller on this seam uses, and closing is the owner's job — {@code Fleetd} + * holds the concrete {@link LeadMailbox} for its shutdown hook and hands only this narrower view to + * everyone else. + */ +public interface LeadChannel { + + /** + * Send {@code m} to {@code toCoordId}'s mailbox, blocking until the broker confirms it is + * durably queued. Throws {@link IllegalStateException} when it is not — unroutable (nobody owns + * that coord-id), nacked, or unconfirmed within the implementation's timeout. A caller must + * report that as a failed send, never as a delivered one. + */ + void publish(String toCoordId, LeadMessage m); + + /** Non-destructive FIFO snapshot of the messages held for this daemon's own coord-id. */ + List peek(); + + /** Drop {@code msgId} from the held set and ack it on the broker. A no-op if it is not held. */ + void ack(String msgId); + + /** This daemon's own lead coordination id — the mailbox it owns, and the {@code from} it sends as. */ + String selfCoordId(); +} diff --git a/bridged/src/main/java/dev/ltms/fleet/msg/LeadCoordLoop.java b/bridged/src/main/java/dev/ltms/fleet/msg/LeadCoordLoop.java new file mode 100644 index 0000000..ad86a63 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/fleet/msg/LeadCoordLoop.java @@ -0,0 +1,193 @@ +package dev.ltms.fleet.msg; + +import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.herdr.AgentStatus; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.util.List; +import java.util.Map; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.function.Supplier; + +/** + * The receive half of lead-to-lead messaging: a bounded background loop that takes what has arrived + * in this daemon's own {@link LeadChannel} mailbox and types it into the local lead's herdr pane. + * + *

{@code fleet_send{coordId}} is the send half — it publishes to a peer daemon's mailbox and + * returns. Nothing on the receiving side reads that mailbox on its own, because the peer lead is an + * interactive agent, not a service that polls; this loop is what closes the gap. + * + *

Status-gated, exactly like {@link ReplyPushLoop}. A pane may only be injected + * into at a turn boundary ({@link AgentStatus#injectable()} — idle, blocked or done); pasting into + * a live turn corrupts it. So a tick that finds the lead busy simply does nothing and comes back + * later. + * + *

Ack only after delivery. A message is acked — removed from the broker — only + * once {@link AgentControl#send} has actually put it in the pane. Anything not delivered (no lead + * pane resolvable, lead mid-turn, herdr threw) stays unacked and is retried on the next tick, and + * survives a daemon restart because the broker still holds it. The cost of that choice is a + * possible duplicate — the send lands and the ack does not — which is the right way round: a peer + * lead seeing a message twice is a nuisance, a peer lead never seeing it at all is the failure this + * whole path exists to remove. + * + *

One message per tick. The loop delivers at most one held message per tick even + * when several are waiting. Injecting a second one immediately would mean acting on a status read + * taken before the first injection: that first paste starts a turn, and herdr does not + * report the pane as {@code working} the instant it does. Waiting for the next tick means every + * delivery is gated on a status read that already saw the previous one. A backlog therefore drains + * one message per {@code intervalMs}, in FIFO order. + */ +public final class LeadCoordLoop { + + private static final Logger log = LoggerFactory.getLogger(LeadCoordLoop.class); + + /** How an arriving peer message is rendered into the lead's pane — the sender's coord-id, then its text. */ + static final String DELIVERY_FORMAT = "[lead %s] %s"; + + private final LeadChannel channel; + private final AgentControl agents; + private final Supplier> leads; + private final ScheduledExecutorService scheduler; + private final long intervalMs; + + private volatile boolean running; + + /** + * @param channel this daemon's own lead mailbox + * @param agents herdr control, for the status gate and the pane injection + * @param leads live {@code terminal_id → name} view of the leads this daemon recognises — + * read through the supplier on every tick, never snapshotted, so a lead found by + * the tab scan after startup becomes reachable without a restart + * @param scheduler the loop's own scheduler; the caller owns its shutdown + * @param intervalMs how long between ticks + */ + public LeadCoordLoop(LeadChannel channel, AgentControl agents, Supplier> leads, + ScheduledExecutorService scheduler, long intervalMs) { + this.channel = channel; + this.agents = agents; + this.leads = leads; + this.scheduler = scheduler; + this.intervalMs = intervalMs; + } + + /** Begin ticking. Idempotent-ish: calling it twice would schedule two chains, so call it once. */ + public void start() { + running = true; + log.info("lead coordination: delivering peer messages for coord-id {} every {}ms", + channel.selfCoordId(), intervalMs); + scheduleNext(); + } + + /** Stop ticking. In-flight work finishes; nothing further is scheduled. */ + public void close() { + running = false; + } + + private void scheduleNext() { + if (!running) { + return; + } + scheduler.schedule(this::tickAndReschedule, intervalMs, TimeUnit.MILLISECONDS); + } + + private void tickAndReschedule() { + try { + tick(); + } catch (RuntimeException e) { + // Never let one bad tick end the chain — the next one re-reads everything from scratch. + log.warn("lead coordination tick failed: {}", e.toString()); + } + scheduleNext(); + } + + /** + * One tick: deliver at most one held peer message into the local lead's pane and ack it. + * Package-private so a test drives it directly rather than waiting on the scheduler. + */ + void tick() { + List held; + try { + held = channel.peek(); + } catch (RuntimeException e) { + log.debug("lead coordination: cannot read the mailbox this tick: {}", e.toString()); + return; + } + if (held.isEmpty()) { + return; + } + String lead = resolveLocalLead(); + if (lead == null) { + // Left unacked on purpose: the broker keeps holding it until a lead pane exists. + log.debug("lead coordination: {} message(s) waiting but no local lead pane to deliver to", + held.size()); + return; + } + AgentStatus status; + try { + status = agents.status(lead); + } catch (RuntimeException e) { + log.debug("lead coordination: status check failed for lead {}, retrying next tick", lead, e); + return; + } + if (!status.injectable()) { + log.debug("lead coordination: lead {} is {} (not injectable), holding {} message(s)", + lead, status, held.size()); + return; + } + LeadMessage msg = held.getFirst(); + try { + agents.send(lead, DELIVERY_FORMAT.formatted(msg.from(), msg.content())); + } catch (RuntimeException e) { + // Not delivered, so not acked — the broker still has it for the next tick. + log.warn("lead coordination: failed to deliver message {} from {} to lead {}: {}", + msg.msgId(), msg.from(), lead, e.toString()); + return; + } + try { + channel.ack(msg.msgId()); + } catch (RuntimeException e) { + // Delivered but not acked: it will be redelivered, which the javadoc calls out as the + // deliberate direction of this trade. + log.warn("lead coordination: delivered message {} but could not ack it: {}", + msg.msgId(), e.toString()); + return; + } + log.debug("lead coordination: delivered message {} from {} to lead {}", msg.msgId(), msg.from(), lead); + } + + /** + * Which local pane a peer's message is for. The mailbox's {@code selfCoordId} is this daemon's + * one lead identity, so there is exactly one right answer — this only has to find it: + * + *

    + *
  1. a lead whose configured name equals {@code selfCoordId} — the explicit, unambiguous case;
  2. + *
  3. otherwise the sole lead, when this daemon recognises exactly one;
  4. + *
  5. otherwise nothing, and the message waits.
  6. + *
+ * + *

Step 3 is deliberate rather than a guess-the-lead fallback. Picking one of several leads + * arbitrarily would type a peer's message into a pane it was not addressed to, and the message + * would then be acked and gone. Leaving it held costs a delay and nothing else. + */ + private String resolveLocalLead() { + Map known = leads.get(); + if (known.isEmpty()) { + return null; + } + String self = channel.selfCoordId(); + for (var entry : known.entrySet()) { + if (entry.getValue() != null && entry.getValue().equals(self)) { + return entry.getKey(); + } + } + if (known.size() == 1) { + return known.keySet().iterator().next(); + } + log.warn("lead coordination: {} leads are known and none is named \"{}\" — cannot tell which " + + "pane a peer message is for; name one lead after coordinator.selfId to fix this", + known.size(), self); + return null; + } +} diff --git a/bridged/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java b/bridged/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java new file mode 100644 index 0000000..f3c8e68 --- /dev/null +++ b/bridged/src/main/java/dev/ltms/fleet/msg/LeadMailbox.java @@ -0,0 +1,431 @@ +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. + * + *

Single-target, unlike {@link AmqpReplyInbox}. {@code AmqpReplyInbox} + * multiplexes many workers' reply queues under one gateway connection. A {@code LeadMailbox} + * instance is simpler: it owns exactly one queue — this daemon's own + * {@code lead..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. + * + *

Consume-and-hold with deferred manual ack — 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. + * + *

Publishing does not imply owning. {@link #publish} sends to + * {@code lead..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. + * + *

Recovery. 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 LeadChannel, 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 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 pendingBySeq = new ConcurrentSkipListMap<>(); + /** The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery tag. */ + private final ConcurrentHashMap 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 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..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 not 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. + */ + @Override + 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); + } + } + + /** The coord-id whose mailbox this instance owns — the {@code from} of everything it publishes. */ + @Override + public String selfCoordId() { + return selfCoordId; + } + + /** Non-destructive FIFO snapshot of this mailbox's currently-held messages. */ + @Override + public List 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 drain() { + List 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. */ + @Override + 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 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..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()); + } + } +} diff --git a/bridged/src/main/java/dev/ltms/fleet/msg/LeadMessage.java b/bridged/src/main/java/dev/ltms/fleet/msg/LeadMessage.java new file mode 100644 index 0000000..4b7bc3d --- /dev/null +++ b/bridged/src/main/java/dev/ltms/fleet/msg/LeadMessage.java @@ -0,0 +1,22 @@ +package dev.ltms.fleet.msg; + +/** + * Wire envelope for a lead-to-lead message carried over {@link LeadMailbox}. + * + *

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. + * + *

{@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) { +} diff --git a/bridged/src/main/java/dev/ltms/fleet/peer/PeerLauncher.java b/bridged/src/main/java/dev/ltms/fleet/peer/PeerLauncher.java index dd4d583..8b36f66 100644 --- a/bridged/src/main/java/dev/ltms/fleet/peer/PeerLauncher.java +++ b/bridged/src/main/java/dev/ltms/fleet/peer/PeerLauncher.java @@ -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}. + * + *

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 diff --git a/bridged/src/test/java/dev/ltms/fleet/FleetdLeadMailboxSelectionTest.java b/bridged/src/test/java/dev/ltms/fleet/FleetdLeadMailboxSelectionTest.java new file mode 100644 index 0000000..82cc333 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/fleet/FleetdLeadMailboxSelectionTest.java @@ -0,0 +1,144 @@ +package dev.ltms.fleet; + +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 dev.ltms.fleet.config.FleetConfig; +import dev.ltms.fleet.msg.LeadMailbox; +import org.junit.jupiter.api.Test; +import org.slf4j.LoggerFactory; + +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * CB-637: the daemon decides whether lead-to-lead messaging is on in + * {@link Fleetd#openLeadMailbox}, not in the config record — so testing + * {@code Coordinator.isConfigured()} alone would pass even if {@code Fleetd} never honoured it. + * These drive the real selection with an injected env map and an injected opener, so no broker is + * involved and no process environment is mutated. + * + *

The invariant every case shares: the feature turns itself OFF, never takes the daemon down. + * A fleet whose coordination broker is missing, half-configured or unreachable must still start and + * still work exactly as it did before this feature existed. + */ +class FleetdLeadMailboxSelectionTest { + + private static final String SECRET = "c00rdPw"; + private static final String RESOLVED_URI = "amqp://user:" + SECRET + "@coord.example:5672/coord"; + + /** Fake opener: records what it was offered, or fails as an unreachable broker would. */ + private static final class RecordingOpener implements Fleetd.LeadMailboxOpener { + String offeredUri; + String offeredSelfId; + int offeredPrefetch = -1; + boolean unreachable; + + @Override + public LeadMailbox open(String uri, String selfCoordId, int prefetch) { + this.offeredUri = uri; + this.offeredSelfId = selfCoordId; + this.offeredPrefetch = prefetch; + if (unreachable) { + throw new IllegalStateException("cannot connect to AMQP coordination broker at " + uri, + new java.net.ConnectException("Connection refused")); + } + // A real LeadMailbox needs a live connection; nothing here dereferences the result + // beyond a null check, so the "reachable" cases assert on what was OFFERED instead. + return null; + } + } + + private static ListAppender captureFleetdLogs() { + Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class); + ListAppender appender = new ListAppender<>(); + appender.start(); + logger.addAppender(appender); + return appender; + } + + private static String joined(ListAppender appender, Level level) { + return appender.list.stream().filter(e -> e.getLevel() == level) + .map(ILoggingEvent::getFormattedMessage).reduce("", (a, b) -> a + "\n" + b); + } + + @Test + void noCoordinatorBlockLeavesTheFeatureOffSilently() { + var appender = captureFleetdLogs(); + var opener = new RecordingOpener(); + + assertNull(Fleetd.openLeadMailbox(null, Map.of(), opener)); + + assertNull(opener.offeredUri, "nothing configured means nothing is opened"); + assertEquals("", joined(appender, Level.WARN), + "an opt-in feature nobody asked for must not warn on every boot"); + } + + @Test + void opensTheMailboxWhenAUriAndSelfIdAreConfigured() { + var opener = new RecordingOpener(); + var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null); + + Fleetd.openLeadMailbox(coordinator, Map.of(), opener); + + assertEquals(RESOLVED_URI, opener.offeredUri); + assertEquals("mac-opus", opener.offeredSelfId); + assertEquals(LeadMailbox.DEFAULT_PREFETCH, opener.offeredPrefetch, + "an unset prefetch takes the mailbox's own default, not zero"); + } + + @Test + void honoursUriEnvOverALiteralUri() { + var opener = new RecordingOpener(); + var coordinator = new FleetConfig.Coordinator("amqp://stale:stale@old:5672/x", "COORD_URI", + "mac-opus", 8); + + Fleetd.openLeadMailbox(coordinator, Map.of("COORD_URI", RESOLVED_URI), opener); + + assertEquals(RESOLVED_URI, opener.offeredUri, "the secret store wins over clear text"); + assertEquals(8, opener.offeredPrefetch); + } + + @Test + void turnsOffWhenUriEnvDoesNotResolve() { + var opener = new RecordingOpener(); + var coordinator = new FleetConfig.Coordinator(null, "COORD_URI", "mac-opus", null); + + assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener)); + + assertNull(opener.offeredUri); + } + + @Test + void warnsAndStaysOffWhenSelfIdIsMissing() { + var appender = captureFleetdLogs(); + var opener = new RecordingOpener(); + var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, null, null); + + assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener)); + + assertNull(opener.offeredUri, "a mailbox with no owning coord-id has no queue to declare"); + String warns = joined(appender, Level.WARN); + assertTrue(warns.contains("coordinator.selfId"), () -> "say which key is missing: " + warns); + assertFalse(warns.contains(SECRET), () -> "the URI's password must never be logged: " + warns); + } + + @Test + void warnsAndStaysOffWhenTheBrokerIsUnreachableAtBoot() { + var appender = captureFleetdLogs(); + var opener = new RecordingOpener(); + opener.unreachable = true; + var coordinator = new FleetConfig.Coordinator(RESOLVED_URI, null, "mac-opus", null); + + assertNull(Fleetd.openLeadMailbox(coordinator, Map.of(), opener), + "a down coordination broker turns the feature off; it must never take the daemon down"); + + String warns = joined(appender, Level.WARN); + assertTrue(warns.contains("coord.example"), () -> "name the host that failed: " + warns); + assertFalse(warns.contains(SECRET), () -> "with credentials stripped: " + warns); + assertTrue(warns.contains("Connection refused"), () -> "and the real reason: " + warns); + } +} diff --git a/bridged/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java b/bridged/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java index d1cd9dc..d98990d 100644 --- a/bridged/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java +++ b/bridged/src/test/java/dev/ltms/fleet/config/FleetConfigTest.java @@ -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; @@ -42,6 +43,49 @@ class FleetConfigTest { assertTrue(cfg.guard().hostSet().contains("ollama.ltms.dev")); } + @Test + void aProfileWithAnAutoCompactWindowBelowTheAcceptedRangeIsRejectedAtLoad(@TempDir Path dir) + throws Exception { + Path f = dir.resolve("low-window.yaml"); + Files.writeString(f, """ + profiles: + ltms-local: + baseUrl: http://gx00.gw:8000 + autoCompactWindow: 50000 + """); + + IllegalStateException e = assertThrows(IllegalStateException.class, () -> FleetConfig.load(f)); + assertTrue(e.getMessage().contains("ltms-local"), "the offending profile is named"); + assertTrue(e.getMessage().contains("autoCompactWindow")); + } + + @Test + void aProfileWithAnAutoCompactWindowInRangeLoadsFine(@TempDir Path dir) throws Exception { + Path f = dir.resolve("in-range-window.yaml"); + Files.writeString(f, """ + profiles: + ltms-local: + baseUrl: http://gx00.gw:8000 + autoCompactWindow: 250000 + """); + + FleetConfig cfg = FleetConfig.load(f); + assertEquals(250_000, cfg.profiles().get("ltms-local").autoCompactWindow()); + } + + @Test + void aProfileWithNoAutoCompactWindowLeavesItNull(@TempDir Path dir) throws Exception { + Path f = dir.resolve("no-window.yaml"); + Files.writeString(f, """ + profiles: + ltms-local: + baseUrl: http://gx00.gw:8000 + """); + + assertNull(FleetConfig.load(f).profiles().get("ltms-local").autoCompactWindow(), + "unset means off — today's behaviour, unchanged"); + } + @Test void appliesDefaultsForMissingSections(@TempDir Path dir) throws Exception { Path f = dir.resolve("minimal.yaml"); @@ -932,6 +976,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"); diff --git a/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java b/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java new file mode 100644 index 0000000..6a3b66c --- /dev/null +++ b/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpLeadCoordTest.java @@ -0,0 +1,96 @@ +package dev.ltms.fleet.mcp; + +import dev.ltms.fleet.msg.FakeLeadChannel; +import dev.ltms.fleet.msg.LeadMessage; +import io.modelcontextprotocol.spec.McpSchema; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * CB-637: the {@code fleet_send{coordId}} route — a message to a PEER LEAD on another daemon, + * published to its durable mailbox instead of typed into a pane this daemon can reach. + * + *

Hermetic: {@link FakeLeadChannel} replaces the AMQP-backed {@code LeadMailbox}, so these run + * with no broker. What is checked here is routing and refusal — which envelope goes out, and which + * calls are rejected before anything is sent. + */ +class FleetMcpLeadCoordTest { + + private static final String SELF = "mac-opus"; + private static final String PEER = "fleet01-lead"; + + private static String textOf(McpSchema.CallToolResult r) { + return ((McpSchema.TextContent) r.content().getFirst()).text(); + } + + @Test + void publishesAnEnvelopeAddressedFromThisDaemonToThePeer() { + var channel = new FakeLeadChannel(SELF); + + McpSchema.CallToolResult res = + FleetMcp.sendToLead(channel, PEER, "you own the auth layer, I own the config one", null, null); + + assertFalse(res.isError(), () -> "expected a success result, got: " + textOf(res)); + assertEquals(1, channel.published().size()); + LeadMessage sent = channel.published().getFirst(); + assertEquals(SELF, sent.from(), "the sender is this daemon's own coord-id, never an argument"); + assertEquals(PEER, sent.to()); + assertEquals("you own the auth layer, I own the config one", sent.content()); + assertNotNull(sent.msgId()); + assertFalse(sent.msgId().isBlank(), "the envelope carries an idempotency id for dedup on redelivery"); + assertTrue(textOf(res).contains(PEER), "the receipt names the peer it reached"); + } + + @Test + void rejectsCoordIdTogetherWithSessionId() { + var channel = new FakeLeadChannel(SELF); + + McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", "term_worker", null); + + assertTrue(res.isError()); + assertTrue(textOf(res).contains("coordId"), () -> textOf(res)); + assertTrue(textOf(res).contains("sessionId"), () -> "the error must name the conflict: " + textOf(res)); + assertEquals(0, channel.published().size(), "an ambiguous call must send nothing at all"); + } + + @Test + void rejectsCoordIdTogetherWithTurnId() { + var channel = new FakeLeadChannel(SELF); + + McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, "turn_7"); + + assertTrue(res.isError()); + assertTrue(textOf(res).contains("turnId"), () -> textOf(res)); + assertEquals(0, channel.published().size()); + } + + @Test + void saysSoPlainlyWhenLeadCoordinationIsNotConfigured() { + McpSchema.CallToolResult res = FleetMcp.sendToLead(null, PEER, "hi", null, null); + + assertTrue(res.isError()); + assertTrue(textOf(res).contains("lead coordination is not configured"), () -> textOf(res)); + assertTrue(textOf(res).contains("coordinator"), () -> "point at the config block to add: " + textOf(res)); + } + + @Test + void reportsAnUnreachablePeerAsAToolErrorNamingIt() { + var channel = new FakeLeadChannel(SELF) + .failPublishWith("lead message m1 was returned as unroutable (mailbox not owned)"); + + McpSchema.CallToolResult res = FleetMcp.sendToLead(channel, PEER, "hi", null, null); + + assertTrue(res.isError(), "a black-holed message must never be reported as delivered"); + assertTrue(textOf(res).contains(PEER), () -> "the error must name the coordId: " + textOf(res)); + assertTrue(textOf(res).contains("unroutable"), () -> "and carry the broker's reason: " + textOf(res)); + } + + @Test + void requiresContent() { + var channel = new FakeLeadChannel(SELF); + + assertTrue(FleetMcp.sendToLead(channel, PEER, " ", null, null).isError()); + assertEquals(0, channel.published().size()); + } +} diff --git a/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java b/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java index f522887..f99872c 100644 --- a/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java +++ b/bridged/src/test/java/dev/ltms/fleet/mcp/FleetMcpTest.java @@ -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 profiles = Map.of(wcfg.profile(), wcfg); ClaudeCodeLauncher delegate = new ClaudeCodeLauncher( new AgentControl(h), new WorkspaceControl(h), new SubscriptionGuard(Set.of("gx00.gw")), @@ -527,6 +527,35 @@ class FleetMcpTest { assertTrue(out.contains("\"liveStatus\":\"unknown\""), out); } + @Test + void listReportsThisDaemonsOwnCoordIdWhenLeadCoordinationIsOn() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + + McpSchema.CallToolResult res = FleetMcp.listFleet( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null, + FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"), + FleetMcp.QuarantineSource.none(), Map.of(), "", "mac-opus"); + + String out = textOf(res); + // There is no peer-discovery surface yet, so this row answers the one question an operator + // cannot answer any other way: which coord-id a peer must use to reach ME. + assertTrue(out.contains("\"coordinator\""), out); + assertTrue(out.contains("\"selfId\":\"mac-opus\""), out); + } + + @Test + void listOmitsTheCoordinatorRowWhenLeadCoordinationIsOff() { + FakeHerdr h = new FakeHerdr(); + SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw"))); + + McpSchema.CallToolResult res = FleetMcp.listFleet( + workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, Map.of(), ""); + + assertFalse(textOf(res).contains("coordinator"), + "an ordinary fleet's output must be unchanged by this feature"); + } + @Test void capacityUsesThePlacementLiveCount() { FakeHerdr h = new FakeHerdr(); diff --git a/bridged/src/test/java/dev/ltms/fleet/member/ClaudeCodeLauncherTest.java b/bridged/src/test/java/dev/ltms/fleet/member/ClaudeCodeLauncherTest.java index 1659739..175181b 100644 --- a/bridged/src/test/java/dev/ltms/fleet/member/ClaudeCodeLauncherTest.java +++ b/bridged/src/test/java/dev/ltms/fleet/member/ClaudeCodeLauncherTest.java @@ -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 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 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 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 + // /worktrees/ must land the entry in /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 /worktrees/. + 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 /bridged, not . + 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 @@ -969,6 +1138,42 @@ class ClaudeCodeLauncherTest { assertNull(startEnv(herdr).get("ANTHROPIC_MODEL")); } + // ── autoCompactWindow: --autocompact is pinned on the command line, opt-in per profile ─────── + + /** A launcher for a profile identical but for its {@code autoCompactWindow:} — the only variable. */ + private ClaudeCodeLauncher serviceWithAutoCompactWindow(FakeHerdr herdr, Integer window) { + FleetConfig.Profile cfg = profileWithAutoCompactWindow(window); + return new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr), + new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), + _ -> null); + } + + private static FleetConfig.Profile profileWithAutoCompactWindow(Integer window) { + return new FleetConfig.Profile("sonnet", "http://gx00.gw:8000", "claude-sonnet-5", null, + "BRIDGED_WORKER_TOKEN", List.of("ccs", "sonnet"), "tab", "bridged-workers", + "w #{n}", "http://127.0.0.1:8765/mcp", null, null, null, null, null, Map.of(), + null, null, null, null, null, null, null, null, window); + } + + @Test + void aConfiguredAutoCompactWindowIsPassedAsAnAutocompactFlag() { + FakeHerdr herdr = new FakeHerdr(); + serviceWithAutoCompactWindow(herdr, 250_000).spawn("sonnet", null, null); + + List args = spawnedArgs(herdr); + int flag = args.indexOf("--autocompact"); + assertTrue(flag >= 0, "the flag is what survives a wrapper argv like [ccs, sonnet]"); + assertEquals("250000", args.get(flag + 1)); + } + + @Test + void aProfileWithNoAutoCompactWindowGetsNoAutocompactFlag() { + FakeHerdr herdr = new FakeHerdr(); + serviceWithAutoCompactWindow(herdr, null).spawn("sonnet", null, null); + + assertFalse(spawnedArgs(herdr).contains("--autocompact")); + } + // --- CB-539: subscription-profile opt-in ---------------------------------------------------- /** A claude-code profile on the subscription: no baseUrl (by design), no off-sub endpoint. */ diff --git a/bridged/src/test/java/dev/ltms/fleet/member/CompositePeerLauncherTest.java b/bridged/src/test/java/dev/ltms/fleet/member/CompositePeerLauncherTest.java index 98db5df..f8cf234 100644 --- a/bridged/src/test/java/dev/ltms/fleet/member/CompositePeerLauncherTest.java +++ b/bridged/src/test/java/dev/ltms/fleet/member/CompositePeerLauncherTest.java @@ -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); } /** diff --git a/bridged/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java b/bridged/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java index 9643b3c..e0a2c5d 100644 --- a/bridged/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java +++ b/bridged/src/test/java/dev/ltms/fleet/member/OpenCodeLauncherTest.java @@ -494,4 +494,121 @@ class OpenCodeLauncherTest { assertTrue(json.path("provider").isMissingNode(), "without a baseUrl opencode resolves its own provider as before"); } + + // --- autoCompactWindow: opencode has no absolute compact-at-N knob, so this is applied as the + // model's own limit.context, only when model: resolves to "provider/model" ----------------- + + private static FleetConfig.Profile opencodeCfgWithAutoCompactWindow(String model, String baseUrl, + Integer window) { + return new FleetConfig.Profile("gemini", baseUrl, model, null, "BRIDGED_WORKER_TOKEN", + List.of("opencode"), "tab", "bridged-workers", "opencode: {model} #{n}", null, + null, null, null, null, FleetConfig.Profile.KIND_OPENCODE, Map.of(), null, null, + null, null, null, null, null, null, window); + } + + @Test + void autoCompactWindowIsAppliedAsThePerModelContextLimit(@TempDir Path root) throws Exception { + FakeHerdr herdr = new FakeHerdr(); + service(herdr, root, opencodeCfgWithAutoCompactWindow("openai/gpt-5", null, 250_000)).spawn(); + + String cfgPath = startEnv(herdr).get("OPENCODE_CONFIG"); + assertNotNull(cfgPath, "autoCompactWindow alone must trigger config generation, with no MCP" + + " and no custom provider set"); + JsonNode limit = new ObjectMapper().readTree(Path.of(cfgPath).toFile()) + .path("provider").path("openai").path("models").path("gpt-5").path("limit"); + assertEquals(250_000, limit.path("context").asInt()); + assertEquals(16384, limit.path("output").asInt(), + "opencode's limit schema requires both keys; output gets a safe documented default"); + } + + @Test + void autoCompactWindowMergesIntoACustomProviderRatherThanOverwritingIt(@TempDir Path root) + throws Exception { + FakeHerdr herdr = new FakeHerdr(); + service(herdr, root, opencodeCfgWithAutoCompactWindow( + "local-vllm/deepseek-v4-flash", "http://127.0.0.1:8000", 300_000)).spawn(); + + JsonNode provider = new ObjectMapper() + .readTree(Path.of(startEnv(herdr).get("OPENCODE_CONFIG")).toFile()) + .path("provider").path("local-vllm"); + assertEquals("@ai-sdk/openai-compatible", provider.path("npm").asText(), + "addCustomProvider's own fields must survive the later limit merge"); + assertEquals(300_000, provider.path("models").path("deepseek-v4-flash") + .path("limit").path("context").asInt()); + assertEquals("deepseek-v4-flash", provider.path("models").path("deepseek-v4-flash") + .path("name").asText(), + "the model's pre-existing 'name' field must survive the limit merge too"); + } + + @Test + void autoCompactWindowWithNoProviderSlashInModelGetsNoLimitAndAWarn(@TempDir Path root) + throws Exception { + FakeHerdr herdr = new FakeHerdr(); + // mcpUrl set too, only so a config file gets written at all to inspect; a bare model name + // with no other config-triggering knob would leave OPENCODE_CONFIG unset entirely, which is + // also correct (nothing to write) but not what this test is asserting. + service(herdr, root, new FleetConfig.Profile("gemini", null, "some-free-model", null, + "BRIDGED_WORKER_TOKEN", List.of("opencode"), "tab", "bridged-workers", + "opencode: {model} #{n}", "http://127.0.0.1:8765/mcp", null, null, null, null, + FleetConfig.Profile.KIND_OPENCODE, Map.of(), null, null, null, null, null, null, + null, null, 250_000)).spawn(); + + String cfgPath = startEnv(herdr).get("OPENCODE_CONFIG"); + JsonNode json = new ObjectMapper().readTree(Path.of(cfgPath).toFile()); + assertTrue(json.path("provider").isMissingNode(), + "a bare model name cannot be targeted at a specific provider/model limit entry — " + + "no silent no-op, but also no broken partial write"); + } + + // --- 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"); + } } diff --git a/bridged/src/test/java/dev/ltms/fleet/msg/FakeLeadChannel.java b/bridged/src/test/java/dev/ltms/fleet/msg/FakeLeadChannel.java new file mode 100644 index 0000000..38ee2e9 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/fleet/msg/FakeLeadChannel.java @@ -0,0 +1,75 @@ +package dev.ltms.fleet.msg; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +/** + * Hermetic stand-in for {@link LeadChannel}: an in-memory mailbox that records what was published + * and what was acked, with no broker anywhere. + * + *

Its whole point is that {@link LeadMailbox} — the production implementation — owns a live AMQP + * connection, so every test of the code AROUND it would otherwise need a broker and end up tagged + * {@code contract}. The broker round trip is covered once, by {@code LeadMailboxTest}; the routing + * ({@code FleetMcp}) and the delivery ({@link LeadCoordLoop}) are covered here. + * + *

Thread-safe: {@link LeadCoordLoop} calls it from its own scheduler thread while a test reads + * the recorded lists. + */ +public final class FakeLeadChannel implements LeadChannel { + + private final String selfCoordId; + private final List held = Collections.synchronizedList(new ArrayList<>()); + private final List published = Collections.synchronizedList(new ArrayList<>()); + private final List acked = Collections.synchronizedList(new ArrayList<>()); + /** When set, every {@link #publish} throws it — the unroutable/nacked/timed-out peer. */ + private volatile IllegalStateException publishFailure; + + public FakeLeadChannel(String selfCoordId) { + this.selfCoordId = selfCoordId; + } + + /** Make every publish fail as an unreachable peer would. */ + public FakeLeadChannel failPublishWith(String message) { + this.publishFailure = new IllegalStateException(message); + return this; + } + + /** Put a message in this mailbox as if a peer had sent it. */ + public FakeLeadChannel hold(LeadMessage m) { + held.add(m); + return this; + } + + @Override + public void publish(String toCoordId, LeadMessage m) { + if (publishFailure != null) { + throw publishFailure; + } + published.add(m); + } + + @Override + public List peek() { + return List.copyOf(held); + } + + @Override + public void ack(String msgId) { + held.removeIf(m -> m.msgId().equals(msgId)); + acked.add(msgId); + } + + @Override + public String selfCoordId() { + return selfCoordId; + } + + public List published() { + return List.copyOf(published); + } + + public List acked() { + return List.copyOf(acked); + } +} diff --git a/bridged/src/test/java/dev/ltms/fleet/msg/LeadCoordLoopTest.java b/bridged/src/test/java/dev/ltms/fleet/msg/LeadCoordLoopTest.java new file mode 100644 index 0000000..e1150b9 --- /dev/null +++ b/bridged/src/test/java/dev/ltms/fleet/msg/LeadCoordLoopTest.java @@ -0,0 +1,150 @@ +package dev.ltms.fleet.msg; + +import dev.ltms.fleet.herdr.AgentControl; +import dev.ltms.fleet.herdr.FakeHerdr; +import org.junit.jupiter.api.Test; + +import java.util.List; +import java.util.Map; +import java.util.concurrent.ScheduledExecutorService; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * Unit tests for {@link LeadCoordLoop} — the receive half of lead-to-lead messaging. Hermetic + * throughout: a {@link FakeLeadChannel} stands in for the mailbox and {@link FakeHerdr} for the + * pane, so no broker and no herdr daemon is involved. The broker round trip has its own + * {@code contract}-tagged test on {@link LeadMailbox}. + * + *

Every test drives {@link LeadCoordLoop#tick()} directly rather than waiting on the scheduler: + * the schedule itself is one {@code scheduler.schedule} call, while the decisions worth pinning — + * inject or hold, ack or leave unacked — all live in the tick. + */ +class LeadCoordLoopTest { + + private static final String SELF = "mac-opus"; + private static final String PEER = "fleet01-lead"; + private static final String LEAD_TERM = "term_lead"; + + /** No scheduler is needed: nothing here calls start(). */ + private static final ScheduledExecutorService NO_SCHEDULER = null; + + private static LeadCoordLoop loop(LeadChannel channel, FakeHerdr herdr, Map leads) { + return new LeadCoordLoop(channel, new AgentControl(herdr), () -> leads, NO_SCHEDULER, 3_000L); + } + + private static List prompts(FakeHerdr herdr) { + return herdr.calls.stream().filter(c -> c.method().equals("agent.prompt")).toList(); + } + + @Test + void deliversAHeldMessageToTheLeadPaneAndAcksIt() { + var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "the merge is blocked")); + var herdr = new FakeHerdr().agentStatus("idle"); + + loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick(); + + assertEquals(1, prompts(herdr).size(), "an injectable lead must receive the peer's message"); + @SuppressWarnings("unchecked") + Map params = (Map) prompts(herdr).getFirst().params(); + assertEquals(LEAD_TERM, params.get("target"), "delivered to the resolved local lead pane"); + assertEquals("[lead " + PEER + "] the merge is blocked", params.get("text"), + "the sender's coord-id is carried into the pane — the lead must know who to answer"); + assertEquals(List.of("m1"), channel.acked(), "a delivered message is acked off the broker"); + assertTrue(channel.peek().isEmpty(), "and is no longer held"); + } + + @Test + void leavesTheMessageUnackedWhenTheLeadIsMidTurn() { + var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello")); + var herdr = new FakeHerdr().agentStatus("working"); + + loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick(); + + assertEquals(0, prompts(herdr).size(), "never paste into a live turn"); + assertEquals(List.of(), channel.acked(), "an undelivered message must NOT be acked"); + assertEquals(1, channel.peek().size(), "it stays held for the next tick"); + } + + @Test + void leavesTheMessageUnackedWhenNoLeadPaneIsKnown() { + var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello")); + var herdr = new FakeHerdr().agentStatus("idle"); + + loop(channel, herdr, Map.of()).tick(); + + assertEquals(0, prompts(herdr).size()); + assertEquals(List.of(), channel.acked(), "nowhere to deliver is not a reason to drop it"); + assertEquals(1, channel.peek().size()); + } + + @Test + void leavesTheMessageUnackedWhenHerdrRefusesTheInjection() { + var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello")); + var herdr = new FakeHerdr().agentStatus("idle").agentSendFailsWith("agent_not_found"); + + loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick(); + + assertEquals(List.of(), channel.acked(), "a send that threw delivered nothing, so nothing is acked"); + assertEquals(1, channel.peek().size()); + } + + @Test + void resolvesTheLeadByNameWhenSeveralAreKnown() { + var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello")); + var herdr = new FakeHerdr().agentStatus("idle"); + // Two leads on this daemon; only one carries the coord-id the mailbox is owned as. + var leads = new java.util.LinkedHashMap(); + leads.put("term_other", "some-other-lead"); + leads.put(LEAD_TERM, SELF); + + loop(channel, herdr, leads).tick(); + + @SuppressWarnings("unchecked") + Map params = (Map) prompts(herdr).getFirst().params(); + assertEquals(LEAD_TERM, params.get("target"), "the lead named after coordinator.selfId wins"); + } + + @Test + void holdsWhenSeveralLeadsAreKnownAndNoneCarriesTheCoordId() { + var channel = new FakeLeadChannel(SELF).hold(new LeadMessage("m1", PEER, SELF, "hello")); + var herdr = new FakeHerdr().agentStatus("idle"); + var leads = new java.util.LinkedHashMap(); + leads.put("term_one", "lead-one"); + leads.put("term_two", "lead-two"); + + loop(channel, herdr, leads).tick(); + + assertEquals(0, prompts(herdr).size(), + "guessing a pane would type a peer's message into the wrong lead and then ack it"); + assertEquals(List.of(), channel.acked()); + } + + @Test + void deliversOneMessagePerTickSoEachIsGatedOnItsOwnStatusRead() { + var channel = new FakeLeadChannel(SELF) + .hold(new LeadMessage("m1", PEER, SELF, "first")) + .hold(new LeadMessage("m2", PEER, SELF, "second")); + var herdr = new FakeHerdr().agentStatus("idle"); + var loop = loop(channel, herdr, Map.of(LEAD_TERM, SELF)); + + loop.tick(); + assertEquals(List.of("m1"), channel.acked(), "FIFO: the oldest goes first"); + assertEquals(1, prompts(herdr).size(), "the second waits for a fresh status read"); + + loop.tick(); + assertEquals(List.of("m1", "m2"), channel.acked()); + assertEquals(2, prompts(herdr).size()); + assertTrue(channel.peek().isEmpty()); + } + + @Test + void anEmptyMailboxNeverTouchesHerdr() { + var channel = new FakeLeadChannel(SELF); + var herdr = new FakeHerdr().agentStatus("idle"); + + loop(channel, herdr, Map.of(LEAD_TERM, SELF)).tick(); + + assertEquals(0, herdr.calls.size(), "an idle fleet must not poll a pane's status every tick"); + } +} diff --git a/bridged/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java b/bridged/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java new file mode 100644 index 0000000..bc015ac --- /dev/null +++ b/bridged/src/test/java/dev/ltms/fleet/msg/LeadMailboxTest.java @@ -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}. + * + *

Proves the mechanism this ticket adds: a {@link LeadMessage} published to + * {@code lead..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 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 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 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 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..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 awaitPeek(LeadMailbox inbox) throws InterruptedException { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10); + List msgs = inbox.peek(); + while (msgs.isEmpty() && System.nanoTime() < deadline) { + Thread.sleep(50); + msgs = inbox.peek(); + } + return msgs; + } +} diff --git a/bridged/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java b/bridged/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java index be58c7c..5e3778b 100644 --- a/bridged/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java +++ b/bridged/src/test/java/dev/ltms/fleet/rest/FleetAppTest.java @@ -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 profiles = Map.of(wcfg.profile(), wcfg); ClaudeCodeLauncher delegate = new ClaudeCodeLauncher( new AgentControl(herdr), new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),