Compare commits

..

33 Commits

Author SHA1 Message Date
Dai Ha edabccd885 CB-637: relabel lead-comms (was colliding CB-635) + document coordId route in the CLAUDE.md tool table
The lead-to-lead wiring shipped with CB-635 in its comments, but CB-635 is
already the broker.uriEnv / unreachable-broker work. Relabel the mailbox +
fleet_send{coordId} + receive loop to CB-637 so a ticket number names one
feature. Add the cross-host peer-lead row to the primary intent->tool table
(kept byte-identical with the wiki template).
2026-08-24 17:48:51 +02:00
Dai Ha 1dbe3a03fc Merge lead-comms-wiring: wire LeadMailbox into the daemon (fleet_send{coordId} + receive loop) 2026-08-24 17:42:11 +02:00
Dai Ha 6058b8472b lead comms: wire LeadMailbox into the daemon so lead-to-lead messages flow
CI / contract (pull_request) Successful in 41s
CI / build (pull_request) Failing after 1m22s
The sibling ticket landed the mechanism (LeadMailbox, LeadMessage, the
coordinator: config block) but nothing opened it, nothing sent through it, and
nothing read it. This is the wiring.

- LeadChannel: a small interface LeadMailbox now implements (publish/peek/ack
  plus a selfCoordId() accessor). It exists so FleetMcp and the receive loop can
  be tested with a fake instead of a live broker. LeadMailbox's AMQP logic is
  untouched — the diff is the implements clause, four @Override marks and the
  accessor.

- Fleetd.openLeadMailbox: opens this daemon's mailbox after the reply inbox is
  selected, with the same env-injected seam selectReplyInbox uses. Every "off"
  path returns null and the daemon still starts: no coordinator block (silent),
  a uriEnv that does not resolve (INFO), a configured broker with no selfId
  (WARN — a mailbox is named after the coord-id that owns it), or a broker that
  refuses at boot (WARN, credentials stripped). Closed in the ordered shutdown
  hook, after the loop that reads it has stopped.

- fleet_send{coordId}: publishes a LeadMessage(from=selfCoordId, to=coordId) to
  the peer's mailbox and returns the broker-confirmed receipt. coordId is
  mutually exclusive with sessionId/turnId and is rejected by name rather than
  resolved by precedence. An unroutable/nacked/timed-out publish comes back as a
  tool error naming the coordId, never a crash. The worker send/reply path is
  not touched.

- LeadCoordLoop: the receive half. Each tick peeks the mailbox, resolves the
  local lead pane, and — only at a turn boundary — injects "[lead <from>] <text>"
  and acks. Anything not delivered stays unacked and is retried, so a message is
  never dropped; one message per tick, so every delivery is gated on a status
  read that already saw the previous one.

- fleet_list reports {selfId, configured} when coordination is on, so an
  operator can find the coord-id a peer must use to reach them. Omitted
  entirely when it is off.

Tests: 20 new hermetic tests (no broker) across routing, delivery and startup
selection. mvn clean install: Tests run: 944, Failures: 0, Errors: 0, Skipped: 0
— BUILD SUCCESS.
2026-08-24 17:39:29 +02:00
Dai Ha f29968c334 Merge autocompact-window: per-profile autoCompactWindow for claude-code + opencode members 2026-08-24 17:32:53 +02:00
Dai Ha 7b98cca967 LeadMailbox: durable leader-to-leader mailbox over a shared coordination vhost
CI / contract (pull_request) Successful in 44s
CI / build (pull_request) Failing after 1m35s
Adds the broker-side mechanism for lead-to-lead messages across daemons/hosts
(unit 1 of 2): a LeadMessage envelope carrying from/to coord-ids, an
AMQP-backed LeadMailbox modeled closely on AmqpReplyInbox (consume-and-hold,
deferred manual ack, confirm-mode publish, recovery handling), and a new
optional coordinator: config block (separate vhost from broker:, leader
traffic only). Config parsing + accessors only — FleetMcp/Fleetd/Injector/
MessageService and the send path are untouched; wiring is a separate ticket.
2026-08-24 17:23:42 +02:00
Dai Ha 2757bc7185 autoCompactWindow: per-profile bounded auto-compaction, both backends
CI / contract (pull_request) Successful in 41s
CI / build (pull_request) Failing after 1m37s
Add opt-in Integer autoCompactWindow to FleetConfig.Profile (last field,
null/unset = today's behaviour). Validated at config load to [100000,
1000000] — the band Claude Code's own --autocompact flag accepts.

Claude Code: appends --autocompact <window> to argv (mirrors --model),
so it survives the ccs <profile> wrapper.

opencode: has no absolute compact-at-N knob (only compaction.auto/prune/
reserved/tail_turns/preserve_recent_tokens), so the window is applied as
the resolved model's own limit.context (+ a required limit.output:16384
default) in the generated opencode.json, merged via get-or-create nodes
so it does not clobber a custom-provider block. Only applies when model:
resolves to "provider/model"; otherwise logs a WARN naming the profile
rather than silently doing nothing.

Docs added to fleetd.example.yaml explaining the cross-backend semantics
difference (compacts AT the window vs. WITHIN it).
2026-08-24 16:54:40 +02:00
Dai Ha 7655f1b51a CB-634: pin the IDE overlay to the module dir + best-effort auto-open
The overlay pinned project_path to the worktree root. For a repo whose Maven
module is a subdir (this repo's pom is in `bridged/`, not at the root), opening
the root imports no module and every ide_* call resolves nothing. Pin and open
the module dir instead.

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

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

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

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

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

911 tests green.
2026-08-23 16:00:36 +02:00
Dai Ha 3f4ac2b24e CB-635: --check reports whether broker.uriEnv resolves in a login shell
CI / contract (push) Successful in 46s
CI / build (push) Failing after 1m30s
An empty uriEnv no longer stops the daemon (#152), so the failure is quiet: bridged
starts, falls back to the in-memory reply inbox, and held reports stop surviving a
restart. --check is the only thing that says so before the fact. The var name is read
out of bridged.yaml so a renamed key cannot make the check lie.
2026-08-23 14:03:08 +02:00
Dai Ha 4644359128 Merge CB-635: broker.uriEnv keeps the AMQP password out of the config, and an unreachable broker no longer stops the daemon (#151, #152)
CI / contract (push) Successful in 1m11s
CI / build (push) Failing after 1m42s
2026-08-23 13:58:21 +02:00
Dai Ha 95a8dbcea9 broker: uriEnv config + non-fatal unreachable broker on boot
CI / build (pull_request) Failing after 56s
CI / contract (pull_request) Successful in 1m6s
#151: Broker gains uriEnv beside uri, taking the AMQP URI from an env var so
the password stays out of fleetd.yaml (same pattern as auth.tokenEnv). uriEnv
wins when set; isConfigured() treats a uriEnv naming an unset/blank variable
as unconfigured. A configured uriEnv is added to the startup required-secrets
report. Never logs the resolved URI (it carries the password).

#152: AmqpReplyInbox.open throwing at boot no longer stops the daemon. The
selection at the call site catches the failure and falls back to the in-memory
inbox for the process lifetime, warning loudly that durable cross-restart
delivery is off and logging the failed URI with credentials stripped.
2026-08-23 13:49:26 +02:00
ltms 83b50753fe Merge CB-633: constrain a member's environment with a derived allow-list
CI / contract (push) Successful in 1m2s
CI / build (push) Failing after 1m23s
Round 2 fixes the defect that mattered: the scrub lived in .zlogin only, and a
herdr pane is a login shell on macOS but a plain interactive one on Linux. It
would have protected nothing on the vhost it was built for, in silence.

Verified by the lead: 901 tests, 0 failures, mvn clean install green. Both new
tests mutation-checked -- unwiring the control fails one, reverting the scrub to
.zlogin alone fails the other while the login-shell test still passes.

NOT yet deployed: the running daemon still holds the old jar.
2026-08-23 08:22:06 +02:00
Dai Ha 6f968d59a4 CB-633 round 2: the scrub ran on macOS and did nothing on Linux
CI / build (pull_request) Failing after 1m8s
CI / contract (pull_request) Successful in 1m10s
Six review findings, from the peer lead `vms` and a reviewer worker. The first
one is a real defect that would have shipped as a dead control.

1. The scrub only ran in a login shell. It lived in the generated `.zlogin`,
   and zsh reads `.zlogin` only for a login shell. herdr does not open the same
   kind of shell everywhere: measured on herdr 0.8.0, a macOS pane runs `-zsh`
   (login) while a Linux pane runs a plain `/usr/bin/zsh`. So on the vhost this
   was being built for, `.zlogin` never ran and every member kept the whole
   secret store, in silence.

   The scrub body now lives in a generated `scrub.zsh` that BOTH `.zshrc` and
   `.zlogin` source, each after sourcing its own `$HOME` counterpart. Linux
   runs the first, macOS runs both, and the second pass is not merely harmless
   -- it re-scrubs anything the operator's `~/.zlogin` exported after `~/.zshrc`
   had finished. Re-running is idempotent.

2. `INFRASTRUCTURE_PASSTHROUGH` listed names that are not infrastructure:
   ANTHROPIC_AUTH_TOKEN, GITEA_TOKEN, GITEA_HOST, ANTHROPIC_BASE_URL,
   ANTHROPIC_MODEL, CLAUDE_CONFIG_DIR, OPENCODE_CONFIG, BRIDGED_MEMBER. I read
   each injection point and confirmed every one of them reaches `launch.env()`
   only when actually injected, so `allowed.addAll(launch.env().keySet())`
   already covers the legitimate case. As static entries they were pure leak
   surface: a host that happened to export ANTHROPIC_AUTH_TOKEN would have had
   it passed straight through.

3. A missing `scrub-report.txt` at teardown was logged at debug. The report is
   the only evidence the scrub ran at all. Its absence has an innocent reading
   and a serious one, and we cannot tell them apart from the daemon -- so it is
   now a WARN that says exactly that. Logging it at debug is how a control that
   quietly stopped working stays unnoticed.

4. Nothing tested that the control was wired in. Deleting the single
   `applyEnvironmentAllowListPolicy(cfg, launch)` line left all 896 tests green
   while turning the feature completely off -- the CB-586/CB-611 shape again.
   `HerdrPeerLauncherAllowListWiringTest` starts a real spawn and asserts on the
   env that reached herdr. Mutation-checked: unwiring that line fails it.

5. `EnvAllowListScrubTest` now also runs `zsh -i` with no `-l`, which is the
   Linux pane shape, so finding 1 is tested from a Mac. Mutation-checked:
   putting the scrub back in `.zlogin` alone fails that test alone, while the
   login-shell test still passes -- which is exactly the blind spot that let
   the bug through.

6. Two ZDOTDIR leaks closed. A failed spawn has no pane id, so its directory
   was never keyed for teardown; it is now removed on the way out. And
   `deleteOnExit` covers a clean shutdown and nothing else, so `generate` now
   reaps sibling directories older than 24h left by a killed daemon.

Also: `policy:` is lowercased with Locale.ROOT, and the `.zlogin`-only claim is
corrected in fleetd.example.yaml, FleetConfig and HerdrPeerLauncher.

901 tests, 0 failures, `mvn clean install` green.
2026-08-23 08:17:52 +02:00
Dai Ha d432df8e5c CB-633: memberCredentials policy=allow-list — derived ZDOTDIR env scrub
CI / build (pull_request) Failing after 1m4s
CI / contract (pull_request) Successful in 1m19s
Move member environment control out of the pane-creation env overlay
(defeated by any file the login shell sources) into a per-spawn ZDOTDIR
directory whose .zlogin runs LAST, after the operator's whole chain, and
blanks every exported variable not on an allow-list DERIVED from what the
launcher itself injects (profiles' tokenEnv/gitTokenEnv/gitHostEnv/env
keys + an infrastructure set) — never hand-typed.

- memberCredentials.policy: allow-list (deny-by-default/deny-list stay
  default and unchanged); known:/allow: become reporting only under it.
- memberCredentials.sshAuthSock knob, blocked by default; allowing it is
  an explicit decision (operator ssh-agent handle).
- Non-zsh login shell: loud WARN, protection off, fallback to the old
  enumerated-name overlay.
- Scrub writes an 'allowed N of M' denominator report, read at teardown;
  credential-shaped blanked names go to WARN (names only, never values).
- Equality test against a real login zsh from a clean parent: surviving
  non-empty exports EQUAL baseline ∩ derived allow-list.
2026-08-23 07:43:38 +02:00
Dai Ha 4b822731e6 CB-632: rename the member's MCP mount bridge -> fleet
CI / contract (push) Successful in 40s
CI / build (push) Successful in 1m22s
Every launcher writes the bridge's MCP server into the config it hands its peer,
and it named that server "bridge". So a member addressed its tools as
mcp__bridge__fleet_send while the tools themselves are already fleet_*. The
mount is named "fleet" now, and a member's tools are mcp__fleet__*.

The name was a bare literal in three files: ClaudeCodeLauncher and LeadLauncher
build a --mcp-config JSON string, OpenCodeLauncher writes an opencode.json node.
Three hand-written copies of one name is how a rename lands in two of them, so
the name is now one constant, PeerLauncher.MCP_MOUNT_NAME.

The mount name is local to the peer — it is the label its own client puts on the
server, and nothing in the daemon reads it back. Renaming it changes no wire
call.

Tests. Each launcher's test now asserts the mount is named fleet AND that
nothing writes "bridge"; the second half is the part that would have caught a
half-done rename. LeadLauncherTest never checked the name at all, only the URL,
so it gained the assertion rather than had one updated.

CLAUDE.md's role-detection ladder quoted mcp__bridge__* as the marker of a
spawned member. It names mcp__fleet__* now, and says that a member spawned
before this change still reports the old prefix. The portable block stays
byte-identical with the wiki template (wiki 569a917).

Build: cd bridged && mvn clean install, then read target/surefire-reports/*.xml
directly — 884 tests, 0 failures, 0 errors.
2026-08-23 07:04:44 +02:00
Dai Ha e38eac1a33 CB-632: ignore fleetd.yaml too, and fix the docs that named the old example
CI / contract (push) Successful in 48s
CI / build (push) Successful in 1m31s
Unit 3 renamed bridged.example.yaml to fleetd.example.yaml and taught Fleetd to
read fleetd.yaml first. Two things it left behind.

bridged/.gitignore still ignored only bridged.yaml. An operator who follows the
new comment and copies the example to fleetd.yaml gets an untracked live config
holding tokens, and git offers to commit it. Both names are ignored now, because
both names work until the cutover.

Three design docs still pointed readers at bridged.example.yaml, a file that no
longer exists under that name.
2026-08-23 07:00:30 +02:00
Dai Ha 8101290933 CB-632: rename the exported metrics bridged_* -> fleet_*
Part of #145 (CB-632). Doing this NOW, ahead of the rest of the path
renames, for one reason: the vms lead is about to wire a monitoring
dashboard to these series. Renaming a metric after a dashboard points at
it breaks continuity and silently leaves a dead panel. Renaming it before
costs nothing, so it goes first rather than at the cutover.

Nine series renamed, all declared in FleetMetrics.

Two real defects found while doing it:

  - FleetApp had "bridged_auth_failures_total" written as a LITERAL
    instead of using FleetMetrics.AUTH_FAILURES -- a second hand-written
    copy of a name, which is how these drift. It now uses the constant,
    so there is one source for that name again.
  - Nothing guarded the prefix. One test does assert a wire name
    (FleetAppAuthTest checks the real /metrics body for fleet_sessions),
    which is good, but it covers one series out of nine. The literal
    above was covered by nothing at all.

So this adds MetricNamesTest, which reads the constants reflectively
rather than listing them -- a test that lists the nine names is itself a
second hand-written copy, and would pass while a tenth went unchecked.
It asserts its own denominator too: "no name starts with bridged_" is
true of an empty set, so a sweep that found nothing would pass loudly.
Asserting the count of 9 makes a broken sweep fail instead.

Also renamed three herdr contract-test workspace labels, __bridged_* ->
__fleet_*. Those are throwaway workspaces created by ensureWorkspace, not
metrics, but they are the same word.

Verified: mvn clean install green, 52 classes, 881 tests, 0 failures.
Test count is up by 3 -- the new prefix, denominator and uniqueness
checks. No bridged_ string remains anywhere outside wiki/.
2026-08-23 06:59:37 +02:00
Dai Ha 6e7fc12f89 CB-632 unit 3: point text references at fleetd.example.yaml
CI / contract (pull_request) Successful in 1m17s
CI / build (pull_request) Successful in 1m36s
2026-08-23 06:51:48 +02:00
Dai Ha 50df14a50f CB-632 unit 3: prefer fleetd.yaml, fall back to bridged.yaml; log the chosen file 2026-08-23 06:51:35 +02:00
Dai Ha ccd882f2fc CB-632 unit 3: rename the example config file to fleetd.example.yaml 2026-08-23 06:48:16 +02:00
Dai Ha ecc590f344 CB-632 unit 5: rename the daemon and classes in the docs prose
Part of #145 (CB-632). Documentation only, plus one internal literal.

Unit 1 renamed the package and classes, which left every doc describing
classes that no longer exist. This fixes the prose across README.md,
docs/ and bridged/docs/ -- 18 files.

Renamed: dev.ltms.bridged -> dev.ltms.fleet, the five class names, and
"bridged" where it names the daemon as a product rather than a path.

Also renamed two literals, because a doc that disagrees with the code is
worse than one that is out of date:

  - bridged-local-noauth -> fleetd-local-noauth. A placeholder apiKey
    OpenCodeLauncher sends when a profile resolves no token, to a local
    endpoint that does not check it. No test asserts the old string.
  - the vnd.ltms.bridged.* media type in the M4 design doc. It appears
    in no Java file, so nothing implements it yet.

Deliberately NOT renamed, because each is still literally true today and
changes only at the cutover:

  - paths: bridged/, bridged.yaml, bridged.example.yaml, bridged.jar,
    .bridged-worktrees, deploy/dev.ltms.bridged.plist,
    scripts/redeploy-bridged.sh, bridged-launchd-wrapper.sh
  - bridged_* metric names -- renaming these after the monitoring is
    wired would break dashboard continuity, so they move before it is
  - bridge_* MCP tool names, which answer alongside fleet_* on purpose
  - BRIDGED_* environment variables, read by a file outside this repo

Method note: perl, not sed. BSD sed has no \b and no lookaround, and a
word-boundary expression there fails silently. The prose replace uses
(?<![\w./-])bridged(?![\w./-]) so it cannot touch a path or an
identifier, then every remaining hit was read by hand.

Verified: mvn clean install green, 51 classes, 878 tests, 0 failures.
2026-08-23 06:46:34 +02:00
Dai Ha 3b7cdf9365 CB-632: the addendum named a class unit 1 renamed
mcp/BridgeMcp is now mcp/FleetMcp. This line is in the project addendum,
not the canonical block, so the wiki template is untouched -- the sync
check still returns True.

Part of #145.
2026-08-23 06:39:54 +02:00
Dai Ha 9b50dd69d8 CB-632 unit 1: rename the Java package and classes to fleet
Part of #145 (CB-632), under epic #125.

The product is called fleet and the daemon is called fleetd, but the code
still said bridge everywhere. This renames the Java half:

  package dev.ltms.bridged -> dev.ltms.fleet
  Bridged        -> Fleetd          (the main class)
  BridgedConfig  -> FleetConfig
  BridgeMcp      -> FleetMcp
  BridgedApp     -> FleetApp
  BridgedMetrics -> FleetMetrics

The package root is dev.ltms.fleet, not dev.ltms.fleetd. The trailing d
means daemon, which names a process, not a namespace.

What this commit deliberately does NOT change:

  - The module directory stays bridged/, and <finalName> stays bridged.
    The installed launchd plist names bridged/target/bridged.jar and its
    KeepAlive is armed, so renaming the jar on its own strands a restart.
    Both change at the cutover, together with the plist, in one step.
  - The bridge_* MCP tool aliases. CB-622 shipped both names on purpose.
    One test names a local variable viaBridge because it holds the result
    of the deprecated call; the rename collided with it and the compiler
    caught it. That variable is back.
  - BRIDGED_* env var names, and bridged.yaml. Both are operator
    contracts and need a read-both shim, which is a later unit.

Two things a plain search-and-replace would have missed:

  - logback.xml and logback-test.xml name the package twice, once as a
    turboFilter class= attribute. The compiler never checks those.
  - BSD sed does not support \b. The word-boundary expression matched
    nothing and said nothing, while the other ten in the same command
    worked. Checked the leftovers instead of trusting the exit code.

Verified: mvn clean install green, 51 test classes, 878 tests, 0 failures
-- the same count as before the rename.
2026-08-23 06:27:43 +02:00
ltms 08f1a79800 Merge #143: scripts/rename-checkout.sh (CB-624)
CI / contract (push) Successful in 1m13s
CI / build (push) Successful in 1m22s
Adds the rename script for CB-624 (#128), with a read-only --check mode.

Two workers ran this ticket in parallel and were kept apart on purpose: one wrote the script, one inventoried the same rename read-only without seeing it. The inventory found three surfaces the script had missed, all verified on the live machine before being acted on: ~/.claude.json's projects key (trust state, enabledMcpjsonServers, allowedTools, lastSessionId), ~/.config/herdr/session.json, and the JetBrains recentProjects.xml / trusted-paths.xml pair in 2026.1 and 2026.2.

All three are report-only by decision, not by oversight, and the reason is written into the script header: ~/.claude.json is global config written by live sessions, herdr/session.json is live process state, and IntelliJ rewrites its own files when the project is reopened.

While adding them the author found a trap worth recording: JetBrains stores $USER_HOME$/LTMS/claude-bridge, not the absolute path, so a plain absolute-path grep reports 0 hits on files full of them. The script now counts both forms and says which form matched.

Every apply-mode branch is unexecuted. Running it moves the directory this system runs from and stops the daemon the workers talk through, so the script is well-reasoned, not proven. Its first real run is its test. Verified before merge: CI run 215 green on both jobs, bash -n passes. shellcheck is not installed on this machine, so it never ran.
2026-08-23 06:09:18 +02:00
Dai Ha 51bdec22e7 CB-624: report three more rename surfaces as report-only
CI / contract (pull_request) Successful in 55s
CI / build (pull_request) Successful in 1m14s
--check gains section 6 counting ~/.claude.json (the projects entry keyed
by the old absolute path), ~/.config/herdr/session.json, and JetBrains
recentProjects.xml / trusted-paths.xml across every IntelliJIdea* version.
Apply mode's closing summary lists the same three with manual follow-ups.

None of the three is rewritten automatically, on purpose: ~/.claude.json
is global live Claude Code config, herdr session.json is live process
state, and JetBrains rewrites its own files when the project is reopened
at the new path. Header comment states the reason for each so it does not
read as an oversight.

JetBrains stores these paths as its $USER_HOME$ macro rather than a
literal absolute path, so the counter matches both forms and says which
form the hits used.
2026-08-23 06:07:16 +02:00
Dai Ha d43f670285 CB-624: add scripts/rename-checkout.sh to rename the checkout safely
CI / contract (pull_request) Successful in 46s
CI / build (pull_request) Successful in 1m32s
Renames ~/LTMS/claude-bridge -> ~/LTMS/fleetd as one auditable command,
modeled on redeploy-bridged.sh. --check reports every surface holding the
old absolute path (checkout files, launchd plist, Claude Code project
state, worktree .git pointers, running daemon). Apply mode stops the
daemon first (launchctl-aware), moves the checkout and the Claude Code
project-state slug dir derived from both paths, repairs worktree gitdir
pointers, rewrites bridged.yaml and the installed plist if they exist,
restarts, and verifies /healthz plus a fresh 'bridged listening' line
anchored to a pre-stop marker.
2026-08-23 05:55:04 +02:00
ltms c135583402 Merge #139: drop the unused rabbitmq host port mapping
CI / contract (push) Successful in 53s
CI / build (push) Successful in 1m19s
The contract job reaches the broker by network alias (AMQP_URI=amqp://guest:guest@rabbitmq:5672), so the host mapping 5672:5672 was never used. It only bound a port on the runner host, which made two concurrent runs collide and the service container fail to start.

Seen on 2026-08-22: run 209 (PR) contract=success and run 210 (the merge of that same code) contract=failure with every step, including checkout, marked cancelled. Both started at 22:08. Re-running 210 alone on the identical commit passed.

PR #139's own run 212 is green on both jobs with the mapping gone, which proves the alias path still works.

Also the worker-PR proof for CB-623 (#127): opened by the agent account against fleet/fleetd after the org transfer.
2026-08-23 05:41:26 +02:00
Dai Ha 8d2893b67e CB-ci: drop host port mapping for rabbitmq service in contract job
CI / contract (pull_request) Successful in 43s
CI / build (pull_request) Successful in 1m33s
2026-08-23 05:35:58 +02:00
Dai Ha 555715ced9 CB-623: point every path at fleet/fleetd after the org transfer
CI / contract (push) Successful in 44s
CI / build (push) Successful in 1m29s
The repo moved lms/claude-bridge -> fleet/claude-bridge -> fleet/fleetd.
Gitea redirects hold, so most of this is not urgent, but one line was a
real break: the implementer skill posts a worker's PR to a hardcoded
repo path, so every worker PR would have gone to the old address.

  .claude/skills/implementer/SKILL.md  the worker PR endpoint (functional)
  .gitmodules                          wiki submodule URL
  CLAUDE.md + wiki/7-Use-Cases.md      the canonical block, kept byte-identical
  README.md                            clone command and wiki link
  deploy/bridged.service               Documentation=
  plugin/.claude-plugin/plugin.json    homepage + repository
  docs/*.md                            issue and wiki links

The wiki is not a separate repo. /repos/lms/claude-bridge.wiki returns 404
and lms owned no .wiki entity, so the wiki moved with the repo; both the old
and the new wiki SSH URLs resolve to the same sha. Ticket step 4 assumed a
second transfer that does not exist.
2026-08-23 05:22:04 +02:00
ltms 446a9d395c Merge CB-622: canonical block to fleet_*, charter quote fixed
CI / contract (push) Successful in 49s
CI / build (push) Successful in 1m22s
2026-08-22 22:08:44 +02:00
198 changed files with 6483 additions and 1746 deletions
+1 -1
View File
@@ -85,7 +85,7 @@ host (`GITEA_HOST`) into your env for exactly this — the token can create a PR
merge**.
```bash
API="${GITEA_HOST%/}/api/v1/repos/lms/claude-bridge/pulls"
API="${GITEA_HOST%/}/api/v1/repos/fleet/fleetd/pulls"
BRANCH="$(git branch --show-current)"
curl -sS -X POST "$API" \
-H "Authorization: token ${GITEA_TOKEN}" \
-2
View File
@@ -69,8 +69,6 @@ jobs:
env:
RABBITMQ_DEFAULT_USER: guest
RABBITMQ_DEFAULT_PASS: guest
ports:
- 5672:5672
env:
# Service containers are reachable from the job by their network alias on their internal port.
AMQP_URI: amqp://guest:guest@rabbitmq:5672
+1 -1
View File
@@ -1,3 +1,3 @@
[submodule "wiki"]
path = wiki
url = ssh://git@git.ltms.dev:2224/lms/claude-bridge.wiki.git
url = ssh://git@git.ltms.dev:2224/fleet/fleetd.wiki.git
+8 -7
View File
@@ -4,7 +4,7 @@
> **Canonical block.** Everything down to §Layering is the portable bridge charter, copied verbatim
> into every project that mounts the bridge MCP. Keep it byte-identical with the template in the
> wiki ([Use Cases](https://git.ltms.dev/lms/claude-bridge/wiki/7-Use-Cases) → *The portable
> wiki ([Use Cases](https://git.ltms.dev/fleet/fleetd/wiki/7-Use-Cases) → *The portable
> CLAUDE.md block*); improvements go to the template first, then out to each project. Anything
> specific to *this* repo lives under §Project addendum below, never inline above it.
@@ -26,9 +26,9 @@ was bound to. Don't infer what you can ask.
Only if that call is unavailable, fall back to these — each is one-way, so keep reading until one
fires: the reply charter in your system prompt (*"You are a spawned member in the
claude-bridge fleet"*) ⇒ **spawned member**; bridge tools prefixed `mcp__bridge__*` ⇒ **spawned
claude-bridge fleet"*) ⇒ **spawned member**; fleet tools prefixed `mcp__fleet__*` ⇒ **spawned
member** (the launcher fixes that mount name; a primary's mount is named by whoever wrote its
`.mcp.json`, so it varies); `ANTHROPIC_BASE_URL` set ⇒ **spawned member** (Claude-model members run
`.mcp.json`, so it varies — and a member spawned before CB-632 still says `mcp__bridge__*`); `ANTHROPIC_BASE_URL` set ⇒ **spawned member** (Claude-model members run
on a clean env, so its *absence* proves nothing). None of these separate a worker from an architect —
only `fleet_whoami` does. **Still unsure ⇒ act as a worker**, the most restricted member role. The
two mistakes are not symmetric: a primary acting as a worker is refused by the authorization gate —
@@ -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: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
| Message a **peer lead** on this host | `fleet_send{sessionId: <their terminal>, content}` — `fleet_list` → `leads` reports it. Coordination only, **never** a task |
| Message a **peer lead** on another daemon or host | `fleet_send{coordId: <their coord-id>, 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}` |
@@ -185,7 +186,7 @@ must obey belongs in the charter, not here.
## Project addendum — claude-bridge (not part of the canonical block)
- **This repo is the bridge.** The daemon is `bridged`, its MCP mount is `http://127.0.0.1:8765/mcp`,
and the code behind the rules above is `mcp/BridgeMcp` (tools), `auth/Authz` (the role table),
and the code behind the rules above is `mcp/FleetMcp` (tools), `auth/Authz` (the role table),
`mcp/ConnectionIdentity` (connection→role), and `worker/*Launcher` (`REPLY_CHARTER`).
- **Skills available to delegate:** `implementer` (worktree → commit → push → own PR) and
`reviewer` (scoped review → one structured finding). Name one in every delegation.
@@ -271,7 +272,7 @@ Before you call any work done, check the row that matches what you touched:
| a `fleet_*` tool — added, removed, renamed, or its params/semantics | the primary's intent→tool table; any rule that names that tool |
| `Authz` / the role table | invariant 3, and the primary-only vs worker-only claims |
| `ConnectionIdentity` / how a caller is resolved | the `fleet_whoami` paragraph and the fallback ladder |
| `REPLY_CHARTER`, or a launcher's mount/flags | the fallback ladder (`mcp__bridge__*`), and the layering table's top row |
| `REPLY_CHARTER`, or a launcher's mount/flags | the fallback ladder (`mcp__fleet__*`), and the layering table's top row |
| the injector / status gating | invariant 4 |
| worktree provisioning or the parity overlay | the "both roles read this file" premise — it rests on the worker's worktree being a checkout of this repo |
| `.claude/skills/**` | the addendum's skill list, and the "name the playbook" rule |
@@ -287,7 +288,7 @@ is a Roadmap line. A change that touches none of the three earns no entry, and t
outcome rather than an omission.
Then **propagate**: the block in this file and the template in the wiki
([Use Cases](https://git.ltms.dev/lms/claude-bridge/wiki/7-Use-Cases) → *The portable `CLAUDE.md`
([Use Cases](https://git.ltms.dev/fleet/fleetd/wiki/7-Use-Cases) → *The portable `CLAUDE.md`
block*) must stay byte-identical, and other projects carrying the block need the same edit. Verify
rather than trust:
+16 -16
View File
@@ -9,23 +9,23 @@ Sibling of [`crush-bridge`](https://git.ltms.dev/systems/vms) (which drives a he
process*, so it inherits `CLAUDE.md`, hooks, skills, and MCP — just pointed at a
cheaper/local model.
## Leading approach — herdr-centric message server (`bridged`)
## Leading approach — herdr-centric message server (`fleetd`)
A small always-on message server, **`bridged`**, controls
A small always-on message server, **`fleetd`**, controls
[herdr](https://herdr.dev) (an agent multiplexer) over its Unix-socket API and exposes a
clean 2-way messaging API as an **MCP server that both the primary and the workers mount** —
one unified Claude setup and the **sole communication gateway** (REST/SSE stays for non-Claude
clients; any broker is `bridged`-internal, below the gateway).
herdr owns the PTYs, multiplexing, persistence, and **agent-status events**; `bridged` owns
clients; any broker is `fleetd`-internal, below the gateway).
herdr owns the PTYs, multiplexing, persistence, and **agent-status events**; `fleetd` owns
policy (subscription boundary, session lifecycle, status-gated delivery) and the client
contract. A Claude member launches with `ANTHROPIC_BASE_URL` pointed at the gateway,
`https://llm.ltms.dev/anthropic`, plus a bearer token; the lead stays env-clean and calls
`bridged`'s MCP tools. See the wiki's **[13 User Guide](wiki/13-User-Guide.md)** to run it.
`fleetd`'s MCP tools. See the wiki's **[13 User Guide](wiki/13-User-Guide.md)** to run it.
```mermaid
flowchart LR
OPUS["Opus — primary<br/>(Claude Code, env CLEAN)<br/>MCP client"]
subgraph BD["bridged — standalone daemon (not a claude process)"]
subgraph BD["fleetd — standalone daemon (not a claude process)"]
SRV["SERVER face<br/>MCP · REST/SSE · policy"]
CLI["CLIENT face<br/>status-gated injector · herdr socket"]
SRV --> CLI
@@ -47,27 +47,27 @@ flowchart LR
```
- **Subscription boundary:** the *primary* never sets `ANTHROPIC_BASE_URL` (stays on
Pro/Max). Only the *secondary* process is off-subscription — and `bridged` itself is a
Pro/Max). Only the *secondary* process is off-subscription — and `fleetd` itself is a
plain daemon (no Anthropic quota), so it may poll/subscribe freely.
- **One gateway (unified MCP setup):** `bridged` is the **sole communication path** for every
- **One gateway (unified MCP setup):** `fleetd` is the **sole communication path** for every
Claude session. Primary and workers each mount it as an MCP server (one `claude mcp add`
line, same on both) and talk over MCP tools — `fleet_send` / `fleet_reply` /
`fleet_status` (with `fleet_ask` planned for the blocked-worker path). **No Claude session
ever addresses a broker, a peer, or the network
directly**; any queue is `bridged`-internal. MCP tool I/O never sets `ANTHROPIC_BASE_URL`, so
directly**; any queue is `fleetd`-internal. MCP tool I/O never sets `ANTHROPIC_BASE_URL`, so
mounting the bridge is subscription-safe by construction.
**Tool naming:** the tools were renamed from `bridge_*` to `fleet_*` (CB-622). The daemon
still answers the old `bridge_*` names for one release, but they are deprecated — use the
`fleet_*` names.
- **How the primary consumes a reply:** a single **blocking MCP call** (`fleet_send`);
`bridged` holds it open until the worker calls `fleet_reply` or its turn hits
`fleetd` holds it open until the worker calls `fleet_reply` or its turn hits
`agent_status=done`, then returns the reply as the tool result. No cross-turn busy-poll, so
no quota burn. SSE is an optional side-channel for humans/dashboards watching status.
- **Worker → primary** rides `bridged`'s **MCP rendezvous** — the reply resolves the primary's
blocking call (or, for detached work, `bridged` **injects the primary's idle pane** when it's
- **Worker → primary** rides `fleetd`'s **MCP rendezvous** — the reply resolves the primary's
blocking call (or, for detached work, `fleetd` **injects the primary's idle pane** when it's
ready), so *no keystroke-into-primary and no broker are involved, even single-host*. The one
exception: a split-host primary that isn't a herdr pane wakes via its own `Stop`-hook, which
polls **`bridged`** (never a broker). See the wiki for the two topologies.
polls **`fleetd`** (never a broker). See the wiki for the two topologies.
- **Different model per process** sidesteps Claude Code's lack of per-subagent provider
routing — the worker isn't a subagent, it's its own configured process.
- **AgentAPI** ([`coder/agentapi`](https://github.com/coder/agentapi)) is retained only as a
@@ -76,11 +76,11 @@ flowchart LR
## Docs
Full design, setup, and operations live in the **[wiki](https://git.ltms.dev/lms/claude-bridge/wiki)**,
Full design, setup, and operations live in the **[wiki](https://git.ltms.dev/fleet/fleetd/wiki)**,
vendored here as a submodule under [`wiki/`](./wiki):
```bash
git clone --recurse-submodules ssh://git@git.ltms.dev:2224/lms/claude-bridge.git
git clone --recurse-submodules ssh://git@git.ltms.dev:2224/fleet/fleetd.git
# or, after a plain clone:
git submodule update --init
```
@@ -90,7 +90,7 @@ Gitea wiki.
## Status
🟢 **Implemented & dogfooded** — the herdr-centric **`bridged`** message server is built and in
🟢 **Implemented & dogfooded** — the herdr-centric **`fleetd`** message server is built and in
real use: an Opus primary delegates tasks to off-subscription workers that reply through the
bridge (code reviews delegated this way have produced committed bug fixes). Selected as the
primary approach 2026-07-11, superseding the AgentAPI plan (2026-07-08); AgentAPI retained as a
+3 -1
View File
@@ -2,7 +2,9 @@
target/
dependency-reduced-pom.xml
# Local runtime config (copy from bridged.example.yaml)
# Local runtime config (copy from fleetd.example.yaml). Both names: the live file is still
# bridged.yaml until the cutover, and Fleetd reads either one.
fleetd.yaml
bridged.yaml
# CB-505 audit trail + daemon stdout/stderr — runtime records, never source
+2 -2
View File
@@ -57,7 +57,7 @@ Populate it from the **orchestration-side** MCP tools — `fleet_send`, `fleet_s
session in `SessionManager`. That caller is, by construction, the primary. Worker-side tools
(`fleet_reply`, `fleet_ask`) never set it.
- **Config override / pin:** a `primary: { terminal: "<id>" }` block in `BridgedConfig`
- **Config override / pin:** a `primary: { terminal: "<id>" }` block in `FleetConfig`
(nested record, same shape as `Broker`). Lets an operator pin it, or supply it when
derivation can't (see degrade case).
- **Degrade:** if the primary is off-host or in a non-herdr terminal, `terminalForPid`
@@ -106,7 +106,7 @@ in v1; the stop-on-empty loop is sufficient.
resolved terminal is non-null **and not a registered worker session**, seen on an
orchestration-side tool. This never mislabels a worker (workers are in `SessionManager`)
and needs no new env var or argument (identity stays connection-derived, per the existing
`BridgeMcp` invariant).
`FleetMcp` invariant).
2. **Readiness-gate mismatch → dedicated loop.** The existing `Injector` gates delivery on
`ready.test(target)` = `WorkerPresence` (the *worker's* MCP connected). The primary is not
@@ -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 <tokens>` 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
@@ -179,7 +206,7 @@ herdrSocket: ~/.config/herdr/herdr.sock
# bridged now propagates ITS OWN PATH to every worker by default; set `env:`
# only to override that or add more (JAVA_HOME, …). Since the default is the
# daemon's PATH, make sure the daemon is started with a good one — see the PATH
# lines in deploy/dev.ltms.bridged.plist and deploy/bridged.service.
# lines in deploy/dev.ltms.fleet.plist and deploy/bridged.service.
#
# Adapter-owned variables always win over `env:`: ANTHROPIC_BASE_URL and the
# rest of the ANTHROPIC_*/CLAUDE_* wiring are applied after it, so an `env:`
@@ -224,7 +251,7 @@ profiles:
# subscription path no guard would vet the URL, so allowing both would be a way around the
# guard rather than a configuration.
#
# GOTCHA 1 — it is invisible to the startup secret check. `Bridged.reportRequiredSecrets`
# GOTCHA 1 — it is invisible to the startup secret check. `Fleetd.reportRequiredSecrets`
# skips subscription profiles on purpose (they need no token), so a boot log that reports
# every secret as fine says nothing about these profiles.
#
@@ -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
@@ -342,7 +373,7 @@ placement: weighted
# / credentialId. Those are hot because the placement policy (and, for credentialId,
# the CB-578 stage B quarantine check) reads them through a supplier — being config is
# not by itself enough to make a key hot.
# EXCEPT `fleet.leaders`: Bridged.main reads it once at startup to build the lead tab
# EXCEPT `fleet.leaders`: Fleetd.main reads it once at startup to build the lead tab
# scanner and launcher, and neither is rebuilt on reload. A changed/added/removed
# `fleet.leaders` entry is silently accepted — the reload reports "config reloaded"
# with nothing in the deferred list — but has NO effect until you restart. Treat it
@@ -505,20 +536,51 @@ guard:
# through either: the daemon logs a WARN naming any credential-shaped env var it finds on neither
# list (never its value), so a secret added to the store later does not go unnoticed forever.
#
# policy → only "deny-by-default" exists today (an operator-authored deny-list was deliberately
# rejected — see above). An unrecognized value refuses to start, naming it.
# allow → credential names a member legitimately needs. Left OUT of the pane's env overlay
# entirely, so the value the pane's own (login) shell exports passes through untouched.
# known → every credential name the operator's store is known to export. Every name here NOT
# also in `allow` is overlaid with a non-secret sentinel value before the pane's login
# shell runs — real protection only for names the login shell does not itself re-export
# (see the ROUND-2 CORRECTION note above for the ones it does).
# policy → "deny-by-default" (the default; also accepted spelled "deny-list") overlays each
# known-but-not-allowed name BEFORE the pane's login shell runs — real protection only
# where that shell does not re-export the name (see ROUND-2 CORRECTION above). An
# unrecognized value refuses to start, naming it.
# policy → "allow-list" (CB-633) moves the control to a per-spawn ZDOTDIR directory the daemon
# generates and passes through tab.create's env map. Each generated startup file sources
# its ~/ counterpart FIRST and then runs the scrub, so the scrub happens after the
# operator's whole chain and no sourced file can undo it.
# The scrub is sourced from BOTH the generated .zshrc and the generated .zlogin, because
# herdr does not open the same kind of shell everywhere: macOS panes run a LOGIN zsh (so
# .zlogin runs), Linux panes run a plain interactive zsh (so .zlogin never runs at all).
# A scrub in .zlogin alone would be a control that silently does nothing on Linux.
# The allow-list is DERIVED, never typed:
# every profile's tokenEnv/gitTokenEnv/gitHostEnv values and env-map keys, plus an
# infrastructure set (PATH HOME SHELL TERM LANG LC_* TMPDIR USER LOGNAME PWD SHLVL EDITOR
# PAGER JAVA_HOME XDG_* ZDOTDIR), plus whatever keys this spawn's own env overlay carries.
# Adding a profile can therefore only widen the list, never break another spawn's scrub.
# Under this policy `known`/`allow` below become REPORTING ONLY — they feed the gap WARN,
# they are no longer a control. If the member's login shell is NOT zsh, the daemon logs a
# loud WARN saying protection is off and falls back to deny-by-default's overlay.
# Each pane writes a scrub-report.txt naming how many variables it kept of how many it
# saw; the daemon logs that "allowed N of M" line when the pane stops. If the report is
# MISSING the daemon logs a WARN instead — the scrub cannot then be confirmed to have
# run, and a silently-dead control is exactly what this policy exists to prevent.
# allow → credential names a member legitimately needs. Under deny-by-default, left OUT of the
# pane's env overlay entirely, so the value the pane's own (login) shell exports passes
# through untouched. Under allow-list: reporting only.
# known → every credential name the operator's store is known to export. Under deny-by-default,
# every name here NOT also in `allow` is overlaid with a non-secret sentinel value before
# the pane's login shell runs — real protection only for names that shell does not itself
# re-export (see the ROUND-2 CORRECTION note above). Under allow-list: reporting only.
# sshAuthSock → whether SSH_AUTH_SOCK may pass through under allow-list ("allow") or must be
# blanked like any other non-derived name ("block", the default). This is a decision you
# have to make explicitly: SSH_AUTH_SOCK is a handle to YOUR ssh-agent, and a member
# holding it can sign with your keys — it sits in no secret file and looks like no
# credential, which is why it slipped past three earlier tickets (gitea #110). Blocking
# it breaks git over SSH inside members (push/fetch authenticate as you); use HTTPS
# remotes or scoped deploy keys instead of allowing it lightly.
#
# HOT-RELOADABLE the same way `fleet:` is (CB-559): read fresh on every spawn, so editing this list
# and reloading config (or restarting) changes what the NEXT spawn inherits; already-running members
# are unaffected either way.
# memberCredentials:
# policy: deny-by-default
# policy: deny-by-default # or "deny-list", or "allow-list" (CB-633) — see above
# sshAuthSock: block # allow-list only; see the sshAuthSock note above
# allow:
# - AI_GATEWAY_TOKEN # named in a profile's tokenEnv (local/gx) — a member reaching the
# # gateway is by design, not a leak
@@ -593,11 +655,33 @@ guard:
# RabbitMQ speaks the same AMQP 0-9-1, so it is a URI-only swap.
# uri → AMQP connection URI. No trailing slash ⇒ the default vhost "/"; an empty path ("/")
# is vhost "" and will NOT connect. Encode a named vhost as .../%2Fmyvhost.
# uriEnv → CB-151: name of a host env var holding the AMQP URI, preferred over `uri` (wins
# whenever set). The URI carries `user:pass@` inline, so naming a variable keeps the
# password out of fleetd.yaml — same pattern as auth.tokenEnv/Profile.tokenEnv. A
# uriEnv that resolves to an unset or blank variable is treated as NOT configured and
# the daemon falls back to the in-memory inbox, warning loudly.
# prefetch → CB-527: consumer basicQos, capping how many unacked messages the inbox holds
# in-heap per owned target (the rest sits on the broker's durable queue instead of
# growing the JVM heap). Default 32 when omitted.
# broker:
# uri: amqp://guest:guest@127.0.0.1:5672
# uriEnv: LAVINMQ_URI
# prefetch: 32
# Shared cross-host LEADER coordination broker. OMIT this block to leave lead-to-lead messaging
# off entirely (config-only in this ticket — nothing here wires it into a live LeadMailbox yet).
# This is a SEPARATE AMQP vhost from `broker:` above: member/worker inboxes always stay on the
# per-fleet `broker:` vhost, and this vhost carries only leader-to-leader traffic, so two fleets
# whose members must never see each other can still share one coordination vhost for their leads.
# uriEnv → name of a host env var holding the coordination AMQP URI, same convention as
# broker.uriEnv (keeps the credential out of fleetd.yaml). Wins over `uri` when set.
# selfId → this daemon's own lead coord-id — the name its mailbox is owned under
# (lead.<selfId>.inbox), e.g. "mac-opus" or "fleet01-lead". Must be globally unique
# across every daemon sharing this vhost.
# prefetch → consumer basicQos, capping how many unacked messages the mailbox holds in-heap.
# Default 32 when omitted.
# coordinator:
# uriEnv: LEAD_COORD_URI
# selfId: mac-opus
# prefetch: 32
# Active push-to-primary (CB-307 Stage 3). When a worker reply lands with no open fleet_send,
+7 -4
View File
@@ -5,17 +5,17 @@
<modelVersion>4.0.0</modelVersion>
<groupId>dev.ltms</groupId>
<artifactId>bridged</artifactId>
<artifactId>fleetd</artifactId>
<version>1.0.0</version>
<packaging>jar</packaging>
<name>bridged</name>
<description>claude-bridge message server: sole gateway between primary/worker Claude sessions and herdr</description>
<name>fleetd</name>
<description>fleet message server: sole gateway between a lead session, its members, and herdr</description>
<properties>
<maven.compiler.release>25</maven.compiler.release>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<mainClass>dev.ltms.bridged.Bridged</mainClass>
<mainClass>dev.ltms.fleet.Fleetd</mainClass>
<jackson.version>2.19.0</jackson.version>
<javalin.version>6.7.0</javalin.version>
@@ -161,6 +161,9 @@
</dependencies>
<build>
<!-- CB-632: stays 'bridged' until the cutover renames the module dir and the launchd
plist together. The installed plist names bridged/target/bridged.jar and
KeepAlive is armed, so renaming the jar alone strands a restart. -->
<finalName>bridged</finalName>
<plugins>
<plugin>
@@ -1,107 +0,0 @@
package dev.ltms.bridged.peer;
import java.util.List;
import java.util.Set;
/**
* 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
* config; today's single adapter is the {@code ClaudeCodeLauncher} / Claude Code over herdr.
*
* <p>The core delegates spawn and teardown to this interface without knowing how the peer is set
* up. Environment variables, CLI flags, subscription guards, transport (herdr tab/pane) layout,
* and naming conventions are all adapter-private — the core sees only the returned
* {@link PeerHandle} whose {@code id()} is the registry/routing key.
*
* <p>The interface is a superset of what {@code SessionManager} and {@code Bridged.main} call
* on the concrete launcher today.
*/
public interface PeerLauncher {
/**
* 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
* the Claude Code herdr adapter is always {@code MID_TURN_ASK, WORKTREE, ORPHAN_REAP}.
*/
Set<Capability> capabilities();
/**
* The capabilities of the adapter that {@code profileName} resolves to (null/blank → the
* default profile, the same resolution {@link #spawn} uses). Distinct from {@link
* #capabilities()}, which unions every configured adapter: a caller that must know whether
* <em>this</em> profile's backend supports a capability — e.g. {@link Capability#SESSION_RESUME}
* before honoring {@link SpawnRequest#resumeSessionId()} — needs the per-profile answer, not
* the fleet-wide union, or a mixed fleet could OK a resume that lands on a non-supporting
* adapter (CB-584).
*
* @throws IllegalArgumentException if the profile is unknown and no default is configured
*/
Set<Capability> capabilitiesFor(String profileName);
/**
* {@code profileName}/requestedCwd null/blank → default resolution. Returns after the peer
* process is live (env + argv + placement complete). Never returns {@code null}.
*
* @param req the spawn parameters (profile, requested cwd, caller cwd)
* @return a handle whose {@link PeerHandle#id()} is the registry/routing key
* @throws IllegalArgumentException if the profile is unknown and no default is configured
*/
PeerHandle spawn(SpawnRequest req);
/**
* The configured worker profile names — the set of names {@code spawn(profileName)} accepts.
*/
Set<String> profiles();
/**
* The profile a no-argument {@link #spawn(SpawnRequest)} uses, or {@code null} if none is configured.
*/
String defaultProfile();
/**
* Resolve the effective working directory for a spawn {@code req} without actually spawning.
* Resolution order: requestedCwd → profile cwd → callerCwd → daemon cwd.
*
* @return the resolved absolute path, never null/blank
*/
String effectiveCwd(SpawnRequest req);
/**
* The parity-overlay file list for {@code profileName} (default list when unset). Used by
* worktree provisioning to copy config files into the isolated checkout before spawning.
*/
List<String> parityOverlay(String profileName);
/**
* The set of all agents this launcher currently tracks, transport-specific. Each element
* exposes at minimum a pane-like {@code id()} matching this launcher's {@link PeerHandle}
* scheme, plus transport-level status. Callers merge this set with the session registry to
* build a live roster view.
*/
List<?> list();
/**
* Reap orphaned peers left behind by a prior daemon process. Only peers whose naming scheme
* matches this launcher's and whose nonce differs from the current process are eligible.
* Best-effort: a failure to list or to stop any one peer is logged and never aborts startup.
*
* @return the number of orphaned peers reaped
*/
int reapOrphanWorkers();
/**
* Tear a peer down by its registry/routing key ({@link PeerHandle#id()}). Tolerates an
* already-gone peer. Also cleans up launcher-private resources (e.g. empty dedicated tabs)
* when safe to do so.
*/
void stop(String id);
/**
* Discard the context of the peer identified by {@code id}. Implementations must bypass normal
* bridge delivery/turn accounting. Unsupported peer kinds return {@code false} without sending
* a guessed command.
*
* @return {@code true} when a reset was sent and its status transition must settle before reuse
*/
boolean clearContext(String id);
}
@@ -1,56 +1,60 @@
package dev.ltms.bridged;
package dev.ltms.fleet;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.config.ConfigRef;
import dev.ltms.bridged.config.ConfigWatcher;
import dev.ltms.bridged.guard.SubscriptionGuard;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.HerdrClient;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.herdr.LeadTabScanner;
import dev.ltms.bridged.lead.LeadLauncher;
import dev.ltms.bridged.herdr.PaneLocator;
import dev.ltms.bridged.herdr.UnixSocketHerdrClient;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.inject.CompletionResolver;
import dev.ltms.bridged.inject.ExhaustedPatternLookup;
import dev.ltms.bridged.inject.ExhaustionSink;
import dev.ltms.bridged.inject.Injector;
import dev.ltms.bridged.inject.StatusPoller;
import dev.ltms.bridged.inject.TurnListener;
import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.auth.MemberRegistry;
import dev.ltms.bridged.auth.CallerResolver;
import dev.ltms.bridged.mcp.BridgeMcp;
import dev.ltms.bridged.mcp.ConnectionIdentity;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.bridged.health.FleetHealthMonitor;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.mcp.LsofPeerPidLookup;
import dev.ltms.bridged.mcp.LsofProcessCwdLookup;
import dev.ltms.bridged.msg.AmqpReplyInbox;
import dev.ltms.bridged.msg.InMemoryReplyInbox;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.msg.ReplyInbox;
import dev.ltms.bridged.msg.LeadHeartbeatLoop;
import dev.ltms.bridged.msg.ReplyPushLoop;
import dev.ltms.bridged.rest.BridgedApp;
import dev.ltms.bridged.session.GitWorktrees;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.peer.PeerLauncher;
import dev.ltms.bridged.session.SessionReaper;
import dev.ltms.bridged.member.ClaudeCodeLauncher;
import dev.ltms.bridged.member.CompositePeerLauncher;
import dev.ltms.bridged.member.HerdrPeerLauncher;
import dev.ltms.bridged.member.OpenCodeLauncher;
import dev.ltms.bridged.placement.BackendQuarantine;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.ConfigWatcher;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.LeadTabScanner;
import dev.ltms.fleet.lead.LeadLauncher;
import dev.ltms.fleet.herdr.PaneLocator;
import dev.ltms.fleet.herdr.UnixSocketHerdrClient;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
import dev.ltms.fleet.inject.ExhaustionSink;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.StatusPoller;
import dev.ltms.fleet.inject.TurnListener;
import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.auth.MemberRegistry;
import dev.ltms.fleet.auth.CallerResolver;
import dev.ltms.fleet.mcp.FleetMcp;
import dev.ltms.fleet.mcp.ConnectionIdentity;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.health.FleetHealthMonitor;
import dev.ltms.fleet.mcp.PrimaryRegistry;
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;
import dev.ltms.fleet.msg.LeadHeartbeatLoop;
import dev.ltms.fleet.msg.ReplyPushLoop;
import dev.ltms.fleet.rest.FleetApp;
import dev.ltms.fleet.session.GitWorktrees;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.session.SessionReaper;
import dev.ltms.fleet.member.ClaudeCodeLauncher;
import dev.ltms.fleet.member.CompositePeerLauncher;
import dev.ltms.fleet.member.HerdrPeerLauncher;
import dev.ltms.fleet.member.OpenCodeLauncher;
import dev.ltms.fleet.placement.BackendQuarantine;
import io.javalin.Javalin;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.LinkedHashMap;
@@ -59,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;
@@ -72,17 +77,41 @@ import java.util.stream.Collectors;
* starts listening. Before anything else it asserts its own environment is clean —
* {@code bridged} is not a Claude process and must never carry a base_url.
*/
public final class Bridged {
public final class Fleetd {
private static final Logger log = LoggerFactory.getLogger(Bridged.class);
private static final Logger log = LoggerFactory.getLogger(Fleetd.class);
/** 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;
/**
* CB-632: prefer {@code fleetd.yaml} in {@code dir}; fall back to {@code bridged.yaml} when
* the new name is not there. The operator's live file is still named {@code bridged.yaml},
* so the old name keeps working until that file moves.
*/
static Path chooseDefaultConfigFile(Path dir) {
Path fleetd = dir.resolve("fleetd.yaml");
if (Files.exists(fleetd)) {
return fleetd;
}
return dir.resolve("bridged.yaml");
}
static void main(String[] args) {
Path configPath = Path.of(args.length > 0 ? args[0] : "bridged.yaml");
BridgedConfig cfg = BridgedConfig.load(configPath);
Path configPath = args.length > 0 ? Path.of(args[0]) : chooseDefaultConfigFile(Path.of(""));
// CB-632: the config file is being renamed bridged.yaml -> fleetd.yaml. Name the file we
// actually loaded, whichever of the two names it carries.
log.info("Using configuration file {}", configPath);
FleetConfig cfg = FleetConfig.load(configPath);
// CB-594: report which secret env vars the config actually needs, by name, before anything
// else can fail on a silently-empty one. A daemon started without a login shell (launchd)
// boots fine either way — this is the only thing that says so out loud.
@@ -127,8 +156,8 @@ public final class Bridged {
// CB-402: one adapter per configured peer kind, fronted by a composite router. A profile's
// `kind:` selects its adapter — claude-code (the default) and opencode partition the profile
// set — and the composite dispatches each SPI call to the adapter that owns the profile/pane.
Map<String, BridgedConfig.Profile> claudeProfiles = new LinkedHashMap<>();
Map<String, BridgedConfig.Profile> opencodeProfiles = new LinkedHashMap<>();
Map<String, FleetConfig.Profile> claudeProfiles = new LinkedHashMap<>();
Map<String, FleetConfig.Profile> opencodeProfiles = new LinkedHashMap<>();
cfg.profiles().forEach((name, w) -> {
if (w.isOpenCode()) {
opencodeProfiles.put(name, w);
@@ -157,7 +186,7 @@ public final class Bridged {
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
// (checked at spawn) and the exhaustion sink wired in below (written on BACKEND_EXHAUSTED).
// The cooldown is deferred (see BridgedConfig#quarantineCooldownSeconds): it is read once
// The cooldown is deferred (see FleetConfig#quarantineCooldownSeconds): it is read once
// here, at startup, and a config reload only changes it for a daemon restart.
BackendQuarantine quarantine = new BackendQuarantine(System::nanoTime,
TimeUnit.SECONDS.toNanos(cfg.quarantineCooldownSeconds()));
@@ -227,7 +256,7 @@ public final class Bridged {
var leaders = cfg.fleet().leaders();
if (!leaders.isEmpty()) {
Set<String> memberSpaces = cfg.profiles().values().stream()
.map(BridgedConfig.Profile::workspace)
.map(FleetConfig.Profile::workspace)
.filter(Objects::nonNull)
.collect(Collectors.toSet());
Map<String, String> tabToName = new LinkedHashMap<>();
@@ -332,7 +361,7 @@ public final class Bridged {
}
@Override
public void onDelivered(String target, dev.ltms.bridged.msg.TurnToken token) {
public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) {
completion.onDelivered(target, token);
sessions.onDelivered(target, token);
}
@@ -354,18 +383,16 @@ public final class Bridged {
StatusPoller poller = new StatusPoller(agents, injector, Injector.POLL_INTERVAL_MILLIS);
poller.start();
// CB-307: reply inbox. A broker: block (with a uri) selects the AMQP-backed durable adapter;
// absent, bridged stays soft-state on the in-memory inbox. The AMQP inbox owns a broker
// CB-307: reply inbox. A broker: block selects the AMQP-backed durable adapter; absent (or
// 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;
if (cfg.broker() != null && cfg.broker().isConfigured()) {
replyInbox = AmqpReplyInbox.open(cfg.broker().uri(), cfg.broker().prefetchOrDefault());
log.info("reply inbox: AMQP broker (durable) at {} (prefetch={})",
cfg.broker().uri(), cfg.broker().prefetchOrDefault());
} else {
replyInbox = new InMemoryReplyInbox();
log.info("reply inbox: in-memory (soft-state)");
}
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.
@@ -389,7 +416,7 @@ public final class Bridged {
// CB-502: the registry is built before the service and the push loop so send/reply outcomes
// are counted at their single funnel rather than at each of the two caller-facing surfaces.
// CB-512: the push loop takes it too, so nudge outcomes (delivered|exhausted) are counted.
Metrics metrics = BridgedMetrics.create(sessions, replyInbox);
Metrics metrics = FleetMetrics.create(sessions, replyInbox);
var pushLoop = new ReplyPushLoop(primaryRegistry, agents, replyInbox,
pushScheduler, maxReminders, backoffMs, metrics);
// CB-551: idle-lead heartbeat. Opt-in; absent `leadHeartbeat:` this is never constructed, so
@@ -480,21 +507,41 @@ public final class Bridged {
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
}
BridgeMcp mcp = new BridgeMcp(messages, workers, sessions, identity, presence,
primaryRegistry, callers, metrics, new BridgeMcp.CapacitySource(profile -> liveCountRef.get().apply(profile),
FleetMcp mcp = new FleetMcp(messages, workers, sessions, identity, presence,
primaryRegistry, callers, metrics, new FleetMcp.CapacitySource(profile -> liveCountRef.get().apply(profile),
profile -> {
var configured = config.get().profiles().get(profile);
return configured == null ? null : configured.maxLoad();
}, () -> config.get().profiles().keySet(), System::nanoTime),
new BridgeMcp.HealthCoverageSource(() -> {
new FleetMcp.HealthCoverageSource(() -> {
var health = config.get().health();
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
health != null && health.notifications() != null && health.notifications().configured());
}),
new BridgeMcp.QuarantineSource(profile -> {
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.
@@ -515,6 +562,8 @@ public final class Bridged {
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();
@@ -527,10 +576,19 @@ public final class Bridged {
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();
}));
Javalin app = new BridgedApp(herdr, workers, sessions, messages, presence, mcp.servlet(),
Javalin app = new FleetApp(herdr, workers, sessions, messages, presence, mcp.servlet(),
callers, metrics).build();
app.start(cfg.bind().host(), cfg.bind().port());
log.info("bridged listening on {}:{}, herdr socket {}",
@@ -546,7 +604,7 @@ public final class Bridged {
* hazard is a property of spawning. A lead is never spawned: the operator started it and named it
* (or labelled its tab) only once it was up, so there is no boot window to guard.
*
* <p>A lead is also never enrolled in {@link MemberPresence} — {@code BridgeMcp} marks presence
* <p>A lead is also never enrolled in {@link MemberPresence} — {@code FleetMcp} marks presence
* for every spawned member (worker and architect), deliberately, since that map doubles as the
* member roster's availability signal and a lead counted there would show up as an available
* member. So without the second disjunct a lead is permanently un-deliverable: every
@@ -560,22 +618,169 @@ public final class Bridged {
return target -> presence.isPresent(target) || leads.get().containsKey(target);
}
/** Injection seam for {@link #selectReplyInbox}: production binds {@link AmqpReplyInbox#open}. */
@FunctionalInterface
interface AmqpOpener {
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.
*
* <p>Every "off" path returns {@code null}, and each says why at the level it deserves:
*
* <ul>
* <li>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.</li>
* <li>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.</li>
* <li>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.</li>
* <li>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.</li>
* </ul>
*/
static LeadMailbox openLeadMailbox(FleetConfig.Coordinator coordinator, Map<String, String> 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.
* Everything else falls back to the in-memory inbox: no broker block, a blank {@code uri}, a
* {@code uriEnv} whose variable is unset or blank, or a broker unreachable at boot. The two
* lossy paths warn <em>loudly</em> — never silently — because what is lost is durable,
* cross-restart reply delivery. Package-private and env-injected so the selection is testable
* without a real broker or a mutable process environment.
*/
static ReplyInbox selectReplyInbox(FleetConfig.Broker broker, Map<String, String> env, AmqpOpener amqp) {
if (broker == null) {
log.info("reply inbox: in-memory (soft-state)");
return new InMemoryReplyInbox();
}
if (broker.hasUriEnv()) {
// uriEnv is authoritative whenever set (CB-151): the operator moved off clear text, so
// it must not quietly fall back onto a stale literal uri.
if (broker.uri() != null && !broker.uri().isBlank()) {
log.info("broker.uri is ignored because broker.uriEnv={} is set", broker.uriEnv());
}
String effectiveUri = broker.effectiveUri(env);
if (effectiveUri == null) {
log.warn("broker.uriEnv={} is unset or blank — durable AMQP reply inbox DISABLED. "
+ "Replies are soft-state and will not survive a restart. Set {} in the "
+ "daemon's environment (see scripts/redeploy-bridged.sh) and restart to "
+ "use the durable broker inbox.",
broker.uriEnv(), broker.uriEnv());
log.info("reply inbox: in-memory (soft-state)");
return new InMemoryReplyInbox();
}
log.info("reply inbox: AMQP broker (durable) via env var {} (prefetch={})",
broker.uriEnv(), broker.prefetchOrDefault());
return openAmqpOrFallback(effectiveUri, broker.prefetchOrDefault(), "uriEnv " + broker.uriEnv(), amqp);
}
// No uriEnv: the literal uri path (existing behaviour).
if (broker.effectiveUri(env) == null) {
log.info("reply inbox: in-memory (soft-state)");
return new InMemoryReplyInbox();
}
log.info("reply inbox: AMQP broker (durable) (prefetch={})", broker.prefetchOrDefault());
return openAmqpOrFallback(broker.uri(), broker.prefetchOrDefault(), "uri", amqp);
}
/**
* Open the AMQP inbox, falling back to the in-memory inbox for this process lifetime if the
* broker cannot be reached at boot (CB-152). Not silent: the warning says durable delivery is
* off, replies are soft-state and will not survive a restart, plus the source that failed and
* the URI <em>with credentials stripped</em>. Never retries in the background — a broker that
* drops <em>after</em> startup already self-heals via the connection factory's automatic
* recovery; only the boot path is changed here.
*/
private static ReplyInbox openAmqpOrFallback(String effectiveUri, int prefetch, String source,
AmqpOpener amqp) {
try {
return amqp.open(effectiveUri, prefetch);
} catch (IllegalStateException e) {
log.warn("cannot reach AMQP broker ({}, {}) — falling back to the in-memory reply inbox "
+ "for this process lifetime. Durable, cross-restart reply delivery is OFF; "
+ "replies are soft-state and will not survive a restart. Reason: {}",
source, stripCredentials(effectiveUri), reasonOf(e));
log.info("reply inbox: in-memory (soft-state)");
return new InMemoryReplyInbox();
}
}
/** An AMQP URI carries {@code user:pass@} inline — show the host/port, never the credentials. */
static String stripCredentials(String uri) {
return uri == null ? null : uri.replaceAll("://[^@/]*@", "://");
}
/** The deepest cause's class and message — the outermost {@code IllegalStateException} echoes the URI (with password). */
private static String reasonOf(Throwable e) {
Throwable t = e;
while (t.getCause() != null && t.getCause() != t) {
t = t.getCause();
}
String msg = t.getMessage();
return t.getClass().getSimpleName() + (msg == null || msg.isBlank() ? "" : ": " + msg);
}
/**
* CB-594: which env vars the loaded config actually needs, and why — every non-{@code
* subscription} profile's {@code tokenEnv} (a subscription profile never reads one, see
* {@link BridgedConfig.Profile#isSubscription()}), plus every profile's {@code gitTokenEnv}
* where set (opt-in). Derived from the config, not hard-coded, so a new profile is covered for
* free. A var required by more than one profile is one entry naming every profile that needs
* it. Deliberately excludes {@code auth.tokenEnv}: that one is already enforced loudly, by a
* startup throw in {@code main()} — about 370 lines <em>below</em> this method's call site
* ({@link #reportRequiredSecrets(BridgedConfig)}), not a few lines above it. That throw only
* {@link FleetConfig.Profile#isSubscription()}), plus every profile's {@code gitTokenEnv}
* where set (opt-in), plus a configured {@code broker.uriEnv} (CB-151). Derived from the
* config, not hard-coded, so a new profile is covered for free. A var required by more than one
* profile is one entry naming every profile that needs it. Deliberately excludes {@code
* auth.tokenEnv}: that one is already enforced loudly, by a startup throw in {@code main()} —
* about 370 lines <em>below</em> this method's call site
* ({@link #reportRequiredSecrets(FleetConfig)}), not a few lines above it. That throw only
* fires when {@code auth.mode: token} is configured; under the default loopback-trust mode it
* never runs, and {@code auth.tokenEnv} is simply not required.
*
* <p>Package-private and pure (no I/O, no logging) so the derivation is unit-testable without
* capturing log output; {@link #reportRequiredSecrets(BridgedConfig)} is the logging caller.
* capturing log output; {@link #reportRequiredSecrets(FleetConfig)} is the logging caller.
*/
static Map<String, List<String>> requiredSecretEnvVars(BridgedConfig cfg) {
static Map<String, List<String>> requiredSecretEnvVars(FleetConfig cfg) {
Map<String, List<String>> requiredBy = new LinkedHashMap<>();
cfg.profiles().forEach((name, profile) -> {
if (!profile.isSubscription()) {
@@ -587,6 +792,11 @@ public final class Bridged {
.add("profile '" + name + "' gitTokenEnv");
}
});
FleetConfig.Broker broker = cfg.broker();
if (broker != null && broker.hasUriEnv()) {
requiredBy.computeIfAbsent(broker.uriEnv(), _ -> new ArrayList<>())
.add("broker uriEnv");
}
return requiredBy;
}
@@ -599,7 +809,7 @@ public final class Bridged {
* <p>A missing entry only warns — it must never refuse to start. A daemon that boots and says
* what is wrong is strictly more useful than one that will not boot at all.
*/
private static void reportRequiredSecrets(BridgedConfig cfg) {
private static void reportRequiredSecrets(FleetConfig cfg) {
Map<String, List<String>> requiredBy = requiredSecretEnvVars(cfg);
if (requiredBy.isEmpty()) {
log.info("startup secrets: no profile references a token env var — nothing to check");
@@ -622,7 +832,7 @@ public final class Bridged {
/**
* CB-596: {@code known:} empty (block absent entirely, or present but empty) means {@link
* BridgedConfig.MemberCredentials#blockedSet()} is empty too — every member pane inherits the
* FleetConfig.MemberCredentials#blockedSet()} is empty too — every member pane inherits the
* operator's whole secret store, unblocked, exactly the defect this ticket fixes. Unlike a
* missing token ({@link #reportRequiredSecrets}), there is no name to point at: the point is
* that the block itself is missing. Warn once at startup and say what to add; never refuse to
@@ -632,17 +842,21 @@ public final class Bridged {
* <p>Package-private so the test can capture the log directly, the same way {@link
* #requiredSecretEnvVars} is exposed for {@link #reportRequiredSecrets}'s own test.
*/
static void reportMemberCredentialsGap(BridgedConfig cfg) {
BridgedConfig.MemberCredentials creds = cfg.memberCredentials();
static void reportMemberCredentialsGap(FleetConfig cfg) {
FleetConfig.MemberCredentials creds = cfg.memberCredentials();
if (creds != null && !creds.known().isEmpty()) {
log.info("memberCredentials: {} known name(s), {} allowed — blocking {} on every spawn",
creds.known().size(), creds.allow().size(), creds.blockedSet().size());
log.info("memberCredentials: policy={}, {} known name(s), {} allowed — blocking {} on "
+ "every spawn{}",
creds.policy(), creds.known().size(), creds.allow().size(), creds.blockedSet().size(),
creds.isAllowList()
? " (allow-list: known/allow are reporting only — the control is the derived ZDOTDIR scrub)"
: "");
return;
}
log.warn("memberCredentials: absent or empty — the daemon will start anyway, and every "
+ "member pane inherits the operator's WHOLE secret store, unblocked (CB-592's "
+ "protection is lost). Add a memberCredentials: block (policy/allow/known) to "
+ "bridged.yaml — see bridged.example.yaml — and restart.");
+ "bridged.yaml — see fleetd.example.yaml — and restart.");
}
/**
@@ -678,6 +892,6 @@ public final class Bridged {
}
}
private Bridged() {
private Fleetd() {
}
}
@@ -1,4 +1,4 @@
package dev.ltms.bridged.auth;
package dev.ltms.fleet.auth;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -1,9 +1,9 @@
package dev.ltms.bridged.auth;
package dev.ltms.fleet.auth;
/**
* The authorization table (CB-505), stated once and enforced on both entry paths.
*
* <p>Most of these rules are already true de facto — {@code BridgeMcp} derives a worker's identity
* <p>Most of these rules are already true de facto — {@code FleetMcp} derives a worker's identity
* from the connection rather than reading it from an argument, so a worker has never been able to
* reply <em>as</em> another worker over MCP. What was missing is that the REST surface trusted the
* session id in the URL path, and neither surface checked role at all. This class makes the
@@ -1,7 +1,7 @@
package dev.ltms.bridged.auth;
package dev.ltms.fleet.auth;
import dev.ltms.bridged.mcp.ConnectionIdentity;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.fleet.mcp.ConnectionIdentity;
import dev.ltms.fleet.peer.MemberRole;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
@@ -12,7 +12,7 @@ import java.util.function.Supplier;
/**
* Resolves every caller to a {@link Principal}, for both entry paths into the core (CB-501).
*
* <p>There are two of them and they are not layered the way the docs suggest: {@code BridgeMcp}
* <p>There are two of them and they are not layered the way the docs suggest: {@code FleetMcp}
* calls the service layer directly and is mounted as a raw servlet (so it never passes through a
* Javalin filter), while the REST routes historically resolved no identity at all. Both now
* delegate here, so the authorization rules are stated once instead of drifting apart.
@@ -59,7 +59,7 @@ public final class CallerResolver {
* <p>Like {@link #leadTerminals}, a supplier rather than a fixed map, so a binding injected
* after startup — when the later spawn lifecycle establishes a live architect session, or an
* operator pins one — takes effect without a restart. Consulted per resolve; today's wiring
* in {@code Bridged} reads a constant from config, which is the degenerate live case.
* in {@code Fleetd} reads a constant from config, which is the degenerate live case.
*/
private final Supplier<Map<String, String>> architectTerminals;
private final Function<String, MemberRole> memberSlotRoles;
@@ -235,7 +235,7 @@ public final class CallerResolver {
// loopback-trust: same-host callers that are not workers are the primary. A non-loopback
// caller is anonymous even here — and startup refuses that combination anyway
// (BridgedConfig.validateAuthExposure), so this is defence in depth, not the control.
// (FleetConfig.validateAuthExposure), so this is defence in depth, not the control.
return isLoopback(remoteAddr) ? Principal.primary(c.pid()) : Principal.anonymous();
}
@@ -1,6 +1,6 @@
package dev.ltms.bridged.auth;
package dev.ltms.fleet.auth;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.fleet.peer.MemberRole;
/** Optional session lifecycle hook for live member-slot bindings. */
public interface MemberLifecycle {
@@ -1,7 +1,7 @@
package dev.ltms.bridged.auth;
package dev.ltms.fleet.auth;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.peer.MemberRole;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -59,7 +59,7 @@ public final class MemberRegistry implements MemberLifecycle {
private final Map<String, String> terminalToSlot = new HashMap<>();
/** Flatten every role pool in {@code fleet} into one registry. Leaders are not members. */
public MemberRegistry(BridgedConfig.Fleet fleet) {
public MemberRegistry(FleetConfig.Fleet fleet) {
Map<String, Entry> flat = new LinkedHashMap<>();
if (fleet != null) {
for (MemberRole role : MemberRole.values()) {
@@ -1,4 +1,4 @@
package dev.ltms.bridged.auth;
package dev.ltms.fleet.auth;
/**
* A resolved caller: its {@link Role}, and — for a worker — the herdr {@code terminal_id} that
@@ -1,4 +1,4 @@
package dev.ltms.bridged.auth;
package dev.ltms.fleet.auth;
/**
* What a caller is allowed to be on the bus (CB-501).
@@ -1,4 +1,4 @@
package dev.ltms.bridged.config;
package dev.ltms.fleet.config;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -16,7 +16,7 @@ import java.util.function.Supplier;
/**
* The daemon's live configuration, re-readable without a restart (CB-559).
*
* <p>Consumers hold this, not a {@link BridgedConfig}, and read through {@link #get()} at the point
* <p>Consumers hold this, not a {@link FleetConfig}, and read through {@link #get()} at the point
* of use. A component that captures {@code ref.get()} into a field at construction has opted out of
* reload — which is sometimes right (see <em>deferred</em> below), but it must then be a deliberate
* choice rather than an accident of where the field was initialised.
@@ -31,7 +31,7 @@ import java.util.function.Supplier;
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Those
* three are read through a supplier on {@code CompositePeerLauncher}, which is what makes
* them hot — not the fact that they are config. <strong>This does NOT include
* {@code fleet.leaders}</strong>: {@code Bridged.main} reads {@code cfg.fleet().leaders()}
* {@code fleet.leaders}</strong>: {@code Fleetd.main} reads {@code cfg.fleet().leaders()}
* once at startup to build the {@code LeadTabScanner} and the {@code LeadLauncher}, and
* neither is reconstructed on reload — so a lead added, removed, or re-{@code tab}'d under
* {@code fleet.leaders} needs a restart, the same as any deferred key below.</li>
@@ -42,7 +42,7 @@ import java.util.function.Supplier;
* {@code guard:}, {@code worktreeRoot:}, adding or removing a profile (a new backend needs its own launcher,
* which is constructed once), <em>and an existing profile's launch settings</em> —
* {@code model}, {@code baseUrl}, {@code argv}, {@code env}, {@code mcpUrl},
* {@code exhaustedPattern} (CB-578 stage A — compiled once into {@code Bridged.main}'s
* {@code exhaustedPattern} (CB-578 stage A — compiled once into {@code Fleetd.main}'s
* pattern map at startup), and the rest. {@code credentialId} (CB-578 stage B) is NOT on
* this list — it is read live off the config supplier at every quarantine check and
* exhaustion event, exactly like {@code weight} / {@code maxLoad}, so it is hot instead.
@@ -65,7 +65,7 @@ import java.util.function.Supplier;
* running. A config file being edited is normally read once mid-save; degrading a working daemon
* because it caught a half-written file would be a bad trade.
*/
public final class ConfigRef implements Supplier<BridgedConfig> {
public final class ConfigRef implements Supplier<FleetConfig> {
private static final Logger log = LoggerFactory.getLogger(ConfigRef.class);
@@ -74,21 +74,21 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
Set.of("bind", "herdrSocket", "broker", "auth");
private final Path path;
private final AtomicReference<BridgedConfig> current;
private final AtomicReference<FleetConfig> current;
public ConfigRef(Path path, BridgedConfig initial) {
public ConfigRef(Path path, FleetConfig initial) {
this.path = path;
this.current = new AtomicReference<>(Objects.requireNonNull(initial, "initial config"));
}
/** A fixed reference that never reloads — for tests and for wiring built from a config in code. */
public static ConfigRef fixed(BridgedConfig cfg) {
public static ConfigRef fixed(FleetConfig cfg) {
return new ConfigRef(null, cfg);
}
/** The live configuration. Read this per use; do not cache it in a field. */
@Override
public BridgedConfig get() {
public FleetConfig get() {
return current.get();
}
@@ -149,10 +149,10 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
if (path == null) {
return Outcome.failed("this config was built in code and has no file to reload from");
}
BridgedConfig old = current.get();
BridgedConfig fresh;
FleetConfig old = current.get();
FleetConfig fresh;
try {
fresh = BridgedConfig.load(path);
fresh = FleetConfig.load(path);
// The same gate startup runs. A config that would have refused to boot must not be able
// to slip in through a reload — that is how a daemon ends up in a state it could never
// have started in, which is the hardest kind to debug.
@@ -182,7 +182,7 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
}
/** Cold keys whose value differs between the running config and the candidate. */
private static List<String> changedColdKeys(BridgedConfig old, BridgedConfig fresh) {
private static List<String> changedColdKeys(FleetConfig old, FleetConfig fresh) {
List<String> changed = new ArrayList<>();
if (!Objects.equals(old.bind(), fresh.bind())) {
changed.add("bind");
@@ -202,7 +202,7 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
}
/** Changed keys that were accepted but whose effect waits for a restart. */
private static List<String> changedDeferredKeys(BridgedConfig old, BridgedConfig fresh) {
private static List<String> changedDeferredKeys(FleetConfig old, FleetConfig fresh) {
List<String> changed = new ArrayList<>();
if (!Objects.equals(old.lifecycle(), fresh.lifecycle())) {
changed.add("lifecycle");
@@ -226,9 +226,9 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
if (!Objects.equals(old.quarantineCooldownSeconds(), fresh.quarantineCooldownSeconds())) {
changed.add("quarantineCooldownSeconds");
}
Map<String, BridgedConfig.Profile> before =
Map<String, FleetConfig.Profile> before =
old.profiles() == null ? Map.of() : old.profiles();
Map<String, BridgedConfig.Profile> after =
Map<String, FleetConfig.Profile> after =
fresh.profiles() == null ? Map.of() : fresh.profiles();
// Adding or removing a profile is deferred: a new backend needs its own launcher, and
// launchers are built once at startup.
@@ -247,7 +247,7 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
// a reload can produce, because the operator has no reason to doubt it.
List<String> relaunch = new ArrayList<>();
before.forEach((name, was) -> {
BridgedConfig.Profile now = after.get(name);
FleetConfig.Profile now = after.get(name);
if (now != null && !sameLaunchSettings(was, now)) {
relaunch.add(name);
}
@@ -266,7 +266,7 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
* {@code CompositePeerLauncher}/the CB-578 stage B exhaustion sink) and really do take effect on
* the next spawn.
*/
private static boolean sameLaunchSettings(BridgedConfig.Profile a, BridgedConfig.Profile b) {
private static boolean sameLaunchSettings(FleetConfig.Profile a, FleetConfig.Profile b) {
return Objects.equals(a.baseUrl(), b.baseUrl())
&& Objects.equals(a.model(), b.model())
&& Objects.equals(a.configDir(), b.configDir())
@@ -276,6 +276,9 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
&& 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())
@@ -283,7 +286,7 @@ public final class ConfigRef implements Supplier<BridgedConfig> {
&& Objects.equals(a.kind(), b.kind())
&& Objects.equals(a.env(), b.env())
&& Objects.equals(a.subscription(), b.subscription())
// CB-578 stage B: exhaustedPattern is compiled once into Bridged.main's pattern map
// CB-578 stage B: exhaustedPattern is compiled once into Fleetd.main's pattern map
// at startup (see ExhaustedPatternLookup wiring) — a reload never re-reads it, so a
// changed pattern must be reported as deferred, exactly like model/baseUrl/argv.
&& Objects.equals(a.exhaustedPattern(), b.exhaustedPattern());
@@ -1,4 +1,4 @@
package dev.ltms.bridged.config;
package dev.ltms.fleet.config;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -1,13 +1,14 @@
package dev.ltms.bridged.config;
package dev.ltms.fleet.config;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.core.JsonParser;
import com.fasterxml.jackson.core.JsonToken;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
import dev.ltms.bridged.msg.AmqpReplyInbox;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.placement.PlacementPolicies;
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;
import org.slf4j.LoggerFactory;
@@ -26,7 +27,7 @@ import java.util.Set;
/**
* {@code bridged} configuration, loaded from a YAML file (see
* {@code bridged.example.yaml}). Unknown keys are ignored so config can grow ahead
* {@code fleetd.example.yaml}). Unknown keys are ignored so config can grow ahead
* of the code — but an unknown <em>top-level</em> key is logged as a WARN at load (CB-530), because
* silently dropping a whole block is indistinguishable from honouring it.
*
@@ -69,9 +70,14 @@ 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 BridgedConfig(
public record FleetConfig(
Bind bind,
String herdrSocket,
Map<String, Profile> profiles,
@@ -89,10 +95,23 @@ public record BridgedConfig(
Auth auth,
ConfigReload configReload,
Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials) {
MemberCredentials memberCredentials,
Coordinator coordinator) {
/** Back-compat form before the {@code coordinator:} block was added. */
public FleetConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
ConfigReload configReload, Integer quarantineCooldownSeconds,
MemberCredentials memberCredentials) {
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
configReload, quarantineCooldownSeconds, memberCredentials, null);
}
/** Back-compat form before the CB-596 {@code memberCredentials:} block was added. */
public BridgedConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
public FleetConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
@@ -106,7 +125,7 @@ public record BridgedConfig(
public static final int DEFAULT_QUARANTINE_COOLDOWN_SECONDS = 1800;
/** Back-compat 14-arg form — no {@code configReload:} block, so file watching stays off. */
public BridgedConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
public FleetConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, String placement, Auth auth) {
@@ -115,7 +134,7 @@ public record BridgedConfig(
}
/** Back-compat form before the optional {@code health:} block was added. */
public BridgedConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
public FleetConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, String placement, Auth auth, ConfigReload configReload) {
@@ -124,7 +143,7 @@ public record BridgedConfig(
}
/** Back-compat form before the CB-578 stage B {@code quarantineCooldownSeconds} field was added. */
public BridgedConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
public FleetConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
@@ -147,7 +166,7 @@ public record BridgedConfig(
* {@code WeightedRoundRobinPolicy}), so losing it makes equal-weight placement
* non-reproducible across restarts. Unmodifiable-wrap instead of copy-and-scramble.
*/
public BridgedConfig {
public FleetConfig {
if (profiles != null && !profiles.isEmpty()) {
Map<String, Profile> normalized = new LinkedHashMap<>();
profiles.forEach((name, p) -> normalized.put(name,
@@ -242,7 +261,7 @@ public record BridgedConfig(
* instead, naming the profile and the key. Live means any session the
* registry still owns (acquired and not yet released), in any state.
* @param kind which peer launcher spawns this profile: {@code "claude-code"} (default —
* the {@link dev.ltms.bridged.member.ClaudeCodeLauncher}) or {@code "opencode"}.
* the {@link dev.ltms.fleet.member.ClaudeCodeLauncher}) or {@code "opencode"}.
* The {@code CompositePeerLauncher} routes {@code spawn}/reap by this value, so
* each adapter drives only its own kind. Normalised to lower-case; blank ⇒ the
* default. It selects the adapter, not the transport — placement, tabs, cwd, and
@@ -275,6 +294,19 @@ public record BridgedConfig(
* 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
* <tokens>} 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
* <em>within</em>, 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,9 +321,13 @@ public record BridgedConfig(
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.bridged.member.ClaudeCodeLauncher} (the default). */
/** Peer kind spawned by {@link dev.ltms.fleet.member.ClaudeCodeLauncher} (the default). */
public static final String KIND_CLAUDE_CODE = "claude-code";
/** Peer kind spawned by the opencode adapter (CB-402). */
public static final String KIND_OPENCODE = "opencode";
@@ -348,6 +384,21 @@ public record BridgedConfig(
// "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 BridgedConfig(
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<String> argv,
String placement, String workspace, String tabLabel, String mcpUrl,
String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv,
String kind, Map<String, String> env, Float weight, Integer maxLoad,
Boolean subscription, String exhaustedPattern, String credentialId,
String ideMcpUrl, 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<String> argv,
String placement, String workspace, String tabLabel, String mcpUrl,
String cwd, List<String> parityOverlay, String gitTokenEnv, String gitHostEnv,
String kind, Map<String, String> env, Float weight, Integer maxLoad,
Boolean subscription, String exhaustedPattern, String credentialId, String ideMcpUrl) {
this(profile, baseUrl, model, configDir, tokenEnv, argv, placement, workspace, tabLabel,
mcpUrl, cwd, parityOverlay, gitTokenEnv, gitHostEnv, kind, env, weight, maxLoad,
subscription, exhaustedPattern, credentialId, ideMcpUrl, null, null);
}
/** True when this profile is served by the Claude Code adapter (the default kind). */
@@ -441,7 +532,7 @@ public record BridgedConfig(
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. */
@@ -472,7 +563,7 @@ public record BridgedConfig(
* {@code ANTHROPIC_AUTH_TOKEN} is injected from {@code tokenEnv}. An {@code env:} entry for
* either is a bypass vector — it is what the subscription path would otherwise let survive
* unguarded — so it is rejected at config load (see
* {@link BridgedConfig#validateSubscriptionProfiles()}).
* {@link FleetConfig#validateSubscriptionProfiles()}).
*/
public boolean envCarriesAnthropicBinding() {
return env != null
@@ -489,6 +580,19 @@ public record BridgedConfig(
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.
@@ -551,16 +655,52 @@ public record BridgedConfig(
*
* @param uri AMQP connection URI, e.g. {@code amqp://guest:guest@127.0.0.1:5672/}. Blank/
* {@code null} ⇒ the broker block is treated as absent (in-memory adapter).
* Ignored when {@code uriEnv} is set.
* @param uriEnv name of a host env var holding the AMQP URI (CB-151). The URI carries its
* credentials inline, so giving the variable <em>name</em> keeps the password
* out of the config file, same as {@code auth.tokenEnv}/{@code Profile.tokenEnv}.
* Wins over {@code uri} whenever set. Blank/{@code null} ⇒ ignored.
* @param prefetch CB-527: the consumer's {@code basicQos} prefetch count, bounding how many
* unacked messages the AMQP inbox holds in-heap per owned target. {@code null}/
* non-positive ⇒ {@link AmqpReplyInbox#DEFAULT_PREFETCH}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Broker(String uri, Integer prefetch) {
public record Broker(String uri, String uriEnv, Integer prefetch) {
/** True when a usable broker URI is configured (an empty block does not enable AMQP). */
/** 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 broker URI is configured (an empty block does not enable AMQP).
* Honors {@code uriEnv} first: if it names a variable that is unset or blank, the broker is
* <em>not</em> configured (the daemon falls back to the in-memory inbox) — a bare {@code uri}
* is only consulted when no {@code uriEnv} is set. Env lookup makes this process-dependent;
* callers already reading {@link System#getenv} are the right ones to invoke it.
*/
public boolean isConfigured() {
return uri != null && !uri.isBlank();
return effectiveUri() != null;
}
/**
* The effective AMQP URI to connect with. {@code uriEnv} wins when set (both over {@code uri}
* and alone). When {@code uriEnv} names a variable that is unset or blank, returns {@code
* null} rather than falling back to {@code uri} — an operator who moved to the secret store
* must not silently drop back onto a stale clear-text URI. Returns the literal {@code uri}
* when no {@code uriEnv} is configured.
*/
public String effectiveUri() {
return effectiveUri(System.getenv());
}
/** As {@link #effectiveUri()}, reading the variable value from {@code env} (the injection seam). */
public String effectiveUri(Map<String, String> env) {
if (hasUriEnv()) {
String value = env.get(uriEnv);
return (value != null && !value.isBlank()) ? value : null;
}
return (uri != null && !uri.isBlank()) ? uri : null;
}
/** The prefetch to use, defaulting to {@link AmqpReplyInbox#DEFAULT_PREFETCH} when unset. */
@@ -569,6 +709,82 @@ public record BridgedConfig(
}
}
/**
* Shared cross-host leader coordination broker: a {@link LeadMailbox} lets two leads on
* different daemons — possibly different hosts — exchange durable messages, which a herdr pane
* injection (how {@code fleet_send} reaches a lead today) cannot do at all. Its mere presence is
* config only in this ticket: nothing here opens a live {@code LeadMailbox} yet, that wiring is
* a separate ticket.
*
* <p>Deliberately a SEPARATE vhost from {@link Broker}, not a reuse of it. {@link Broker} is
* per-fleet — its queues are named by worker session id, and two fleets sharing one broker stay
* isolated by vhost (see {@code Two fleets share one LavinMQ}). Leader coordination is meant to
* cross exactly that boundary: two independently-owned fleets' leads talking to each other. Using
* the same vhost would either leak member traffic across the fleet boundary this is meant to
* cross, or force every fleet's members onto one shared vhost to get leader coordination — a
* second, dedicated vhost keeps "member inboxes stay fleet-local" true while still letting leads
* reach across fleets.
*
* @param uri AMQP connection URI for the coordination vhost, e.g.
* {@code amqp://guest:guest@127.0.0.1:5672/coord}. Blank/{@code null} ⇒ the
* coordinator block is treated as unconfigured. Ignored when {@code uriEnv} is set.
* @param uriEnv name of a host env var holding the AMQP URI, same convention as
* {@link Broker#uriEnv()} — keeps the credential out of the config file. Wins
* over {@code uri} whenever set. Blank/{@code null} ⇒ ignored.
* @param selfId this daemon's own lead coord-id — the name its {@code LeadMailbox} is owned
* under ({@code lead.<selfId>.inbox}), e.g. {@code "mac-opus"}. Blank/
* {@code null} ⇒ kept as {@code null} (no self id configured).
* @param prefetch the consumer's {@code basicQos} prefetch count. {@code null}/non-positive ⇒
* {@link LeadMailbox#DEFAULT_PREFETCH}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Coordinator(String uri, String uriEnv, String selfId, Integer prefetch) {
public Coordinator {
selfId = (selfId == null || selfId.isBlank()) ? null : selfId;
}
/** True when a {@code uriEnv} is configured by name, whether or not its variable resolves. */
public boolean hasUriEnv() {
return uriEnv != null && !uriEnv.isBlank();
}
/**
* True when a usable coordination broker URI is configured (an empty block does not enable
* it). Honors {@code uriEnv} first: if it names a variable that is unset or blank, the
* coordinator is <em>not</em> configured — a bare {@code uri} is only consulted when no
* {@code uriEnv} is set.
*/
public boolean isConfigured() {
return effectiveUri() != null;
}
/**
* The effective AMQP URI to connect with. {@code uriEnv} wins when set (both over
* {@code uri} and alone). When {@code uriEnv} names a variable that is unset or blank,
* returns {@code null} rather than falling back to {@code uri} — an operator who moved to
* the secret store must not silently drop back onto a stale clear-text URI. Returns the
* literal {@code uri} when no {@code uriEnv} is configured.
*/
public String effectiveUri() {
return effectiveUri(System.getenv());
}
/** As {@link #effectiveUri()}, reading the variable value from {@code env} (the injection seam). */
public String effectiveUri(Map<String, String> env) {
if (hasUriEnv()) {
String value = env.get(uriEnv);
return (value != null && !value.isBlank()) ? value : null;
}
return (uri != null && !uri.isBlank()) ? uri : null;
}
/** The prefetch to use, defaulting to {@link LeadMailbox#DEFAULT_PREFETCH} when unset. */
public int prefetchOrDefault() {
return (prefetch != null && prefetch > 0) ? prefetch : LeadMailbox.DEFAULT_PREFETCH;
}
}
/**
* Optional pinned primary terminal config (CB-307). When present with a non-blank
* {@code terminal}, the bridge uses this as the primary's herdr identity instead of
@@ -893,7 +1109,7 @@ public record BridgedConfig(
*
* <p>Member identity never depends on this block: a loopback peer PID that maps to a herdr
* pane is unforgeable and is always honoured (see
* {@link dev.ltms.bridged.mcp.ConnectionIdentity}). This only decides what happens for
* {@link dev.ltms.fleet.mcp.ConnectionIdentity}). This only decides what happens for
* <em>everyone else</em>.
*
* @param mode {@code "loopback-trust"} (default) — any loopback caller that is not a known
@@ -957,27 +1173,86 @@ public record BridgedConfig(
* UNLESS it is also in {@link #allow}. A name that shows up in neither list is not silently
* allowed — see {@code HerdrPeerLauncher}'s gap detector, which logs it.
*
* @param policy how the block is computed. Only {@link #POLICY_DENY_BY_DEFAULT} is understood
* today; {@code null}/blank defaults to it. An operator's own deny-list is
* deliberately not supported — see above.
* <p><b>CB-633: allow-list.</b> Deny-by-default's overlay is applied BEFORE the pane's login
* shell runs, so any file that chain sources can re-export over it — and did. The allow-list
* policy moves the control to a generated ZDOTDIR whose startup files run the scrub LAST, after
* the whole operator chain, and blank every exported variable not on an allow-list DERIVED from
* what the launcher itself injects (profiles' tokenEnv/gitTokenEnv/gitHostEnv/env keys plus an
* infrastructure set) — never hand-typed, so adding a profile cannot break a spawn. Under this
* policy {@link #allow} and {@link #known} stop being a control and become reporting only.
*
* @param policy how the block is computed. {@link #POLICY_DENY_BY_DEFAULT} (the default; also
* accepted as {@link #POLICY_DENY_LIST}) shadows each {@code known}-but-not-allowed
* name in the pane-creation env overlay — which a login shell that re-exports the
* name defeats (see CB-596's round-2 correction). {@link #POLICY_ALLOW_LIST}
* (CB-633) replaces the overlay with a per-spawn ZDOTDIR scrub that runs AFTER the
* pane's login shell has finished sourcing everything, blanking every variable not
* on the DERIVED allow-list (derived from what the launcher itself injects — never
* hand-typed). Under {@code allow-list}, {@link #allow} and {@link #known} are
* REPORTING ONLY: they feed the gap WARN, they are no longer a control.
* @param allow credential names a member legitimately needs (e.g. the gateway token it reaches
* the LLM through, the repo-scoped forge token it opens its own PR with). Every
* name here is left unmentioned in the pane's env overlay, so the value the pane's
* own (login) shell exports passes through untouched.
* @param known every credential name the operator's store is known to export. Every name here
* that is NOT also in {@link #allow} is overlaid with a non-secret sentinel value,
* shadowing whatever the pane's login shell would otherwise export for it.
* the LLM through, the repo-scoped forge token it opens its own PR with). Under
* deny-list/deny-by-default every name here is left unmentioned in the pane's env
* overlay; under allow-list this list is reporting only — the control is derived,
* not configured.
* @param known every credential name the operator's store is known to export. Under
* deny-list/deny-by-default every name here that is NOT also in {@link #allow} is
* overlaid with a non-secret sentinel value. Under allow-list this list is
* reporting only.
* @param sshAuthSock whether the member may inherit {@code SSH_AUTH_SOCK} under the allow-list
* policy ({@code "allow"}) or must have it blanked ({@code "block"}, the default).
* This is a DECISION, never a default: {@code SSH_AUTH_SOCK} is a handle to the
* operator's ssh-agent, and a member holding it can sign with the operator's own
* keys — but it appears in no secret file and is credential-shaped like nothing on
* any list, which is why three earlier tickets missed it (gitea #110 / CB-607).
* Blocking it breaks git over SSH inside the member; allow it only when members do
* not need to authenticate as the operator over SSH. Ignored under deny-list /
* deny-by-default, which never touch the name.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record MemberCredentials(String policy, List<String> allow, List<String> known) {
public record MemberCredentials(String policy, List<String> allow, List<String> known,
String sshAuthSock) {
/** The only policy this build understands: block every {@code known} name not in {@code allow}. */
/** Default policy: block every {@code known} name not in {@code allow}, via the env overlay. */
public static final String POLICY_DENY_BY_DEFAULT = "deny-by-default";
/** Alias of {@link #POLICY_DENY_BY_DEFAULT}, spelled the way CB-633 names the two policies. */
public static final String POLICY_DENY_LIST = "deny-list";
/**
* CB-633: derive the kept-name set from what the launcher itself injects, generate a
* per-spawn ZDOTDIR whose startup files blank every exported variable not on it AFTER the
* pane's shell has finished sourcing the operator's chain. The scrub is sourced from both
* the generated {@code .zshrc} and {@code .zlogin}, because herdr opens a LOGIN zsh on macOS
* and a plain interactive one on Linux — see {@code EnvAllowListScrub}.
*/
public static final String POLICY_ALLOW_LIST = "allow-list";
/** The pre-CB-633 three-field form — {@code sshAuthSock} defaults to blocked. */
public MemberCredentials(String policy, List<String> allow, List<String> known) {
this(policy, allow, known, null);
}
public MemberCredentials {
policy = (policy == null || policy.isBlank()) ? POLICY_DENY_BY_DEFAULT : policy.toLowerCase();
String normalizedPolicy = (policy == null || policy.isBlank())
? POLICY_DENY_BY_DEFAULT : policy.toLowerCase(java.util.Locale.ROOT);
// deny-list is an alias of deny-by-default, not a third behaviour — normalize to one
// spelling so every isDenyList()-style check has one value to compare against.
policy = POLICY_DENY_LIST.equals(normalizedPolicy) ? POLICY_DENY_BY_DEFAULT : normalizedPolicy;
allow = allow == null ? List.of() : List.copyOf(allow);
known = known == null ? List.of() : List.copyOf(known);
sshAuthSock = (sshAuthSock != null && "allow".equalsIgnoreCase(sshAuthSock.trim()))
? "allow" : "block";
}
/** True when this block selects the CB-633 derived-allow-list policy. */
public boolean isAllowList() {
return POLICY_ALLOW_LIST.equals(policy);
}
/** True when {@code SSH_AUTH_SOCK} may pass through under the allow-list policy. Default: no. */
public boolean sshAuthSockAllowed() {
return "allow".equals(sshAuthSock);
}
/** {@link #allow} as a set, for membership checks. */
@@ -1028,7 +1303,7 @@ public record BridgedConfig(
return defaultProfileFor(MemberRole.DEV);
}
private static final Logger log = LoggerFactory.getLogger(BridgedConfig.class);
private static final Logger log = LoggerFactory.getLogger(FleetConfig.class);
private static final ObjectMapper YAML = new ObjectMapper(new YAMLFactory());
@@ -1037,17 +1312,17 @@ public record BridgedConfig(
* {@link #warnUnknownTopLevelKeys}. Keep in step with the record components.
*
* <p>Package-private (not {@code private}) so a test can assert every key here is documented in
* {@code bridged.example.yaml} — the only committed description of the config schema, since
* {@code fleetd.example.yaml} — the only committed description of the config schema, since
* {@code bridged.yaml} itself is gitignored.
*/
static final Set<String> KNOWN_TOP_LEVEL_KEYS = Set.of(
"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 BridgedConfig load(Path path) {
public static FleetConfig load(Path path) {
try {
String yaml = Files.readString(path);
rejectRenamedTopLevelKeys(yaml);
@@ -1055,11 +1330,12 @@ public record BridgedConfig(
warnUnknownTopLevelKeys(yaml, path);
rejectDuplicateMemberSlots(yaml);
rejectNegativeMaxLoad(yaml);
rejectAutoCompactWindowOutOfRange(yaml);
rejectUnknownKind(yaml);
rejectUnknownAuthMode(yaml);
rejectUnknownPlacement(yaml);
rejectUnknownMemberCredentialsPolicy(yaml);
BridgedConfig cfg = YAML.readValue(yaml, BridgedConfig.class);
FleetConfig cfg = YAML.readValue(yaml, FleetConfig.class);
// CB-606: validated here, eagerly, using PlacementPolicies.fromName as the single source
// of truth — not lazily at first spawn (see CompositePeerLauncher's placementPolicy
// Supplier), where a bad name would still start a daemon that looks healthy.
@@ -1357,6 +1633,52 @@ public record BridgedConfig(
}
}
/** Lowest {@code autoCompactWindow} Claude Code's {@code --autocompact <tokens>} flag accepts. */
static final int AUTO_COMPACT_WINDOW_MIN = 100_000;
/** Highest {@code autoCompactWindow} Claude Code's {@code --autocompact <tokens>} 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 <tokens>} flag accepts (100k–1M), naming both the profile and
* the value.
*
* <p>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<String> 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<String> KNOWN_KINDS = Set.of(Profile.KIND_CLAUDE_CODE, Profile.KIND_OPENCODE);
@@ -1366,7 +1688,7 @@ public record BridgedConfig(
*
* <p>{@link Profile}'s compact constructor only lower-cases {@code kind} and compares it against
* {@code KIND_OPENCODE} — anything else, including a typo like {@code opencod}, silently falls
* into the claude-code bucket ({@link dev.ltms.bridged.member.CompositePeerLauncher} routes by
* into the claude-code bucket ({@link dev.ltms.fleet.member.CompositePeerLauncher} routes by
* exact adapter claim, not by membership in a known set). With {@code argv:} also unset, the argv
* default special-cases only the exact string {@code "claude-code"}, so the launch command falls
* back to {@code List.of(kind)} — the daemon then tries to run a program literally named after the
@@ -1487,9 +1809,15 @@ public record BridgedConfig(
}
}
/** The member-credential policies this build understands — {@link MemberCredentials#policy()}'s only valid value. */
/**
* The member-credential policies this build understands — {@link MemberCredentials#policy()}'s
* only valid values. Checked against the RAW yaml text (before {@link MemberCredentials}'s
* compact constructor normalizes {@code deny-list} onto {@code deny-by-default}), so the alias
* is listed explicitly.
*/
private static final Set<String> KNOWN_MEMBER_CREDENTIALS_POLICIES =
Set.of(MemberCredentials.POLICY_DENY_BY_DEFAULT);
Set.of(MemberCredentials.POLICY_DENY_BY_DEFAULT, MemberCredentials.POLICY_DENY_LIST,
MemberCredentials.POLICY_ALLOW_LIST);
/**
* Reject a {@code memberCredentials.policy} that is not {@link #KNOWN_MEMBER_CREDENTIALS_POLICIES}
@@ -1570,7 +1898,7 @@ public record BridgedConfig(
}
/** Fill in nested defaults so callers never see nulls for structural fields. */
public BridgedConfig withDefaults() {
public FleetConfig withDefaults() {
Bind b = bind != null ? bind : new Bind(null, 0);
Guard g = guard != null ? guard : new Guard(List.of());
Lifecycle l = lifecycle != null ? lifecycle : new Lifecycle(null, null, null, false);
@@ -1601,12 +1929,17 @@ public record BridgedConfig(
// guard/lifecycle, an absent block is not a safe "feature off" default here, it is a gap. It
// is deliberately not pre-populated with a Java-side name list (that would just reintroduce
// the hardcoded-list defect this record replaces); the block must be configured in
// bridged.yaml to protect anything. See bridged.example.yaml's memberCredentials: comment.
// bridged.yaml to protect anything. See fleetd.example.yaml's memberCredentials: comment.
// CB-633: policy stays deny-by-default here — the allow-list scrub is opt-in, because it is
// stricter than today's behaviour (it blanks every non-derived name, not just known ones)
// 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());
return new BridgedConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
// 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);
}
/**
@@ -1639,8 +1972,8 @@ public record BridgedConfig(
* <p>The scan reads a tab label and concludes "a lead lives here". bridged also <em>writes</em>
* tab labels — every member gets one rendered into its tab. Choose a lead {@code tabPrefix} that
* a member template matches and the daemon starts labelling its own members as leads, promoting
* the entire fleet to {@link dev.ltms.bridged.auth.Role#PRIMARY} with no message and no diff.
* The member-space exclusion in {@link dev.ltms.bridged.herdr.LeadTabScanner} already blocks the
* the entire fleet to {@link dev.ltms.fleet.auth.Role#PRIMARY} with no message and no diff.
* The member-space exclusion in {@link dev.ltms.fleet.herdr.LeadTabScanner} already blocks the
* realistic path, but defence that depends on one workspace label holding is not defence enough
* for a privilege boundary.
*
@@ -1,4 +1,4 @@
package dev.ltms.bridged.guard;
package dev.ltms.fleet.guard;
/** Thrown when the subscription boundary would be violated. Never swallow this. */
public class GuardException extends RuntimeException {
@@ -1,4 +1,4 @@
package dev.ltms.bridged.guard;
package dev.ltms.fleet.guard;
import java.net.URI;
import java.net.URISyntaxException;
@@ -1,7 +1,7 @@
package dev.ltms.bridged.health;
package dev.ltms.fleet.health;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.session.MemberSession;
/**
* Pure classifier. Collection and repair are deliberately outside this package.
@@ -1,10 +1,10 @@
package dev.ltms.bridged.health;
package dev.ltms.fleet.health;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.session.MemberSession;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -44,7 +44,7 @@ public final class FleetHealthMonitor {
* invoked once when a member transitions into a terminal health state. Required —
* there is deliberately no defaulting overload; a caller that does not want the
* fail-tickets-on-terminal-health behavior must pass an explicit inert value (see
* {@code TestTurnTokens.inert} / {@code BridgeMcp.CapacitySource.none()} for the pattern).
* {@code TestTurnTokens.inert} / {@code FleetMcp.CapacitySource.none()} for the pattern).
*/
public FleetHealthMonitor(AgentControl agents, Supplier<List<MemberSession>> roster, MessageService messages,
ScheduledExecutorService scheduler, LongSupplier clock, long intervalSeconds,
@@ -1,4 +1,4 @@
package dev.ltms.bridged.health;
package dev.ltms.fleet.health;
/** Classification plus the private fact that the next pure decision needs. */
public record HealthDecision(HealthState state, HealthPrior prior) { }
@@ -1,4 +1,4 @@
package dev.ltms.bridged.health;
package dev.ltms.fleet.health;
/** Private cross-tick observation. It is deliberately not a reported health value. */
public record HealthPrior(boolean busyButDone) {
@@ -1,7 +1,7 @@
package dev.ltms.bridged.health;
package dev.ltms.fleet.health;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.session.MemberSession;
/** Read-only facts from one fleet collection tick. */
public record HealthSnapshot(MemberSession.State sessionState, AgentStatus liveStatus,
@@ -1,4 +1,4 @@
package dev.ltms.bridged.health;
package dev.ltms.fleet.health;
/** Health classifications reported for a member. */
public enum HealthState {
@@ -1,6 +1,6 @@
package dev.ltms.bridged.health;
package dev.ltms.fleet.health;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.fleet.msg.Rendezvous;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.health;
package dev.ltms.fleet.health;
import java.util.ArrayList;
import java.util.HashMap;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.herdr;
package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.herdr;
package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.herdr;
package dev.ltms.fleet.herdr;
/**
* A herdr agent's lifecycle state, as reported by {@code agent_status}. Drives the
@@ -1,4 +1,4 @@
package dev.ltms.bridged.herdr;
package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.herdr;
package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.JsonNode;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.herdr;
package dev.ltms.fleet.herdr;
/**
* Raised when a herdr call fails: transport error, or an {@code error} envelope
@@ -1,4 +1,4 @@
package dev.ltms.bridged.herdr;
package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
import org.slf4j.Logger;
@@ -14,7 +14,7 @@ import java.util.function.Supplier;
/**
* Discovers which panes host a lead by scanning herdr for tabs the operator labelled by convention
* (CB-531), and hands {@link dev.ltms.bridged.auth.CallerResolver} the resulting
* (CB-531), and hands {@link dev.ltms.fleet.auth.CallerResolver} the resulting
* {@code terminal_id → lead name} map.
*
* <p><strong>Why scan at all.</strong> A lead is never spawned — a human opens a tab and starts an
@@ -49,7 +49,7 @@ import java.util.function.Supplier;
* <p><strong>CB-558 — bridged now writes lead labels too.</strong> This class used to be able to say
* that bridged never renames a lead tab, so the label was always the human's own writing and there
* was no round-trip from the daemon's rename back into its next decision.
* {@code dev.ltms.bridged.lead.LeadLauncher} ends that: an auto-launched lead is labelled by the
* {@code dev.ltms.fleet.lead.LeadLauncher} ends that: an auto-launched lead is labelled by the
* daemon and found again by this scan. The trust direction above is unaffected — bridged writing a
* name for a lead it just started is not a pane promoting itself — but <em>staleness</em> becomes
* real: a label left behind by a session that has since died would read as a live lead forever.
@@ -58,7 +58,7 @@ import java.util.function.Supplier;
* ever make a decision that <em>removes</em> something based on this map, add the same check.
* The remaining hazard is an <em>operator</em> one — a worker {@code tabLabel} template that
* happens to start with the same prefix would promote the whole fleet — and that is refused at
* startup by {@code BridgedConfig.validateLeadTabPrefixes} rather than documented here.
* startup by {@code FleetConfig.validateLeadTabPrefixes} rather than documented here.
*
* <p><strong>Caching.</strong> {@link #get()} is on the request path (every resolve), so the scan
* is TTL-cached and a stale-but-valid map is preferred to a herdr round-trip. A failed scan keeps
@@ -1,4 +1,4 @@
package dev.ltms.bridged.herdr;
package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.herdr;
package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.herdr;
package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.herdr;
package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.herdr;
package dev.ltms.fleet.herdr;
import com.fasterxml.jackson.databind.JsonNode;
import org.slf4j.Logger;
@@ -1,8 +1,8 @@
package dev.ltms.bridged.inject;
package dev.ltms.fleet.inject;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.msg.TurnToken;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.TurnToken;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -289,7 +289,7 @@ public final class CompletionResolver implements TurnListener {
/**
* Coverage summary for the CB-578 stage A exhausted-pattern classification, logged at startup
* the way {@link dev.ltms.bridged.health.FleetHealthMonitor#coverage} is — so an operator can
* the way {@link dev.ltms.fleet.health.FleetHealthMonitor#coverage} is — so an operator can
* see whether the classification is on, and for which profiles, without reading every
* profile's config by hand.
*
@@ -1,4 +1,4 @@
package dev.ltms.bridged.inject;
package dev.ltms.fleet.inject;
import java.util.regex.Pattern;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.inject;
package dev.ltms.fleet.inject;
/**
* Notified when {@link CompletionResolver} actually delivers a {@code BACKEND_EXHAUSTED}
@@ -8,7 +8,7 @@ package dev.ltms.bridged.inject;
*
* <p>{@link CompletionResolver} knows only {@code target} (a herdr terminal id); it has no notion of
* profiles or credentials, so mapping {@code target} to whatever should be quarantined is entirely
* the sink's job — see {@code Bridged.main}'s wiring, which resolves target → session → profile →
* the sink's job — see {@code Fleetd.main}'s wiring, which resolves target → session → profile →
* {@code effectiveCredentialId()} and calls {@code BackendQuarantine.quarantine} on it.
*/
@FunctionalInterface
@@ -1,8 +1,8 @@
package dev.ltms.bridged.inject;
package dev.ltms.fleet.inject;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.msg.TurnToken;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.msg.TurnToken;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -81,7 +81,7 @@ public final class Injector {
/**
* The single source for the injector poll cadence — how often the {@link StatusPoller} drives
* {@link #onStatus} at. {@code Bridged} passes this to every {@link StatusPoller} it constructs,
* {@link #onStatus} at. {@code Fleetd} passes this to every {@link StatusPoller} it constructs,
* and this class reads it to state the readiness grace in seconds on the CB-562 expiry log
* instead of hardcoding "60s". One constant, so a cadence change cannot silently desync a log
* that claims a grace duration.
@@ -1,4 +1,4 @@
package dev.ltms.bridged.inject;
package dev.ltms.fleet.inject;
import java.util.concurrent.ConcurrentHashMap;
import java.util.Set;
@@ -1,8 +1,8 @@
package dev.ltms.bridged.inject;
package dev.ltms.fleet.inject;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.HerdrException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -1,7 +1,7 @@
package dev.ltms.bridged.inject;
package dev.ltms.fleet.inject;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -1,6 +1,6 @@
package dev.ltms.bridged.inject;
package dev.ltms.fleet.inject;
import dev.ltms.bridged.msg.TurnToken;
import dev.ltms.fleet.msg.TurnToken;
/**
* Notified when a worker's delegated turn is observed to complete — a confirmed
@@ -1,12 +1,12 @@
package dev.ltms.bridged.lead;
package dev.ltms.fleet.lead;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.herdr.Tab;
import dev.ltms.bridged.herdr.Workspace;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.Tab;
import dev.ltms.fleet.herdr.Workspace;
import dev.ltms.fleet.herdr.WorkspaceControl;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -17,6 +17,7 @@ import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.stream.Collectors;
import dev.ltms.fleet.peer.PeerLauncher;
/**
* Starts the leads {@code fleet.leaders:} declares, when none is already running (CB-558).
@@ -54,14 +55,14 @@ public final class LeadLauncher {
private final AgentControl agents;
private final WorkspaceControl spaces;
private final BridgedConfig cfg;
private final FleetConfig cfg;
/**
* @param agents herdr agent control (start, list)
* @param spaces workspace / tab control (ensure, create, label, list)
* @param cfg the loaded config — {@code fleet.leaders}, {@code profiles} and each lead's tab
*/
public LeadLauncher(AgentControl agents, WorkspaceControl spaces, BridgedConfig cfg) {
public LeadLauncher(AgentControl agents, WorkspaceControl spaces, FleetConfig cfg) {
this.agents = agents;
this.spaces = spaces;
this.cfg = cfg;
@@ -75,7 +76,7 @@ public final class LeadLauncher {
* the remaining leads are still attempted.
*/
public int ensureLeads() {
Map<String, BridgedConfig.Leader> leaders = cfg.fleet().leaders();
Map<String, FleetConfig.Leader> leaders = cfg.fleet().leaders();
if (leaders.isEmpty()) {
return 0;
}
@@ -91,9 +92,9 @@ public final class LeadLauncher {
}
int started = 0;
for (Map.Entry<String, BridgedConfig.Leader> e : leaders.entrySet()) {
for (Map.Entry<String, FleetConfig.Leader> e : leaders.entrySet()) {
String name = e.getKey();
BridgedConfig.Leader lead = e.getValue();
FleetConfig.Leader lead = e.getValue();
int running = live.getOrDefault(name, 0);
int wanted = lead.instances();
@@ -110,7 +111,7 @@ public final class LeadLauncher {
continue;
}
BridgedConfig.Profile profile = cfg.profiles().get(lead.profile());
FleetConfig.Profile profile = cfg.profiles().get(lead.profile());
if (profile == null) {
log.warn("lead '{}' names profile '{}', which is not configured — not launching",
name, lead.profile());
@@ -138,9 +139,9 @@ public final class LeadLauncher {
* already carries {@link Agent#tabId()} directly, so a hand-opened lead is found the same way an
* auto-launched one is — by labelling its tab to match.
*/
private Map<String, Integer> liveLeads(Map<String, BridgedConfig.Leader> leaders) {
private Map<String, Integer> liveLeads(Map<String, FleetConfig.Leader> leaders) {
Set<String> memberSpaces = cfg.profiles().values().stream()
.map(BridgedConfig.Profile::workspace)
.map(FleetConfig.Profile::workspace)
.filter(w -> w != null && !w.isBlank())
.collect(Collectors.toSet());
@@ -174,12 +175,12 @@ public final class LeadLauncher {
* <p>Matched exactly (case-insensitively) against each lead's configured {@code tab}, so an
* operator's {@code "lead: something-else"} tab is not mistaken for a configured lead.
*/
private String leadNameOf(String label, Map<String, BridgedConfig.Leader> leaders) {
private String leadNameOf(String label, Map<String, FleetConfig.Leader> leaders) {
if (label == null) {
return null;
}
String l = label.strip();
for (Map.Entry<String, BridgedConfig.Leader> e : leaders.entrySet()) {
for (Map.Entry<String, FleetConfig.Leader> e : leaders.entrySet()) {
String tab = e.getValue().tabLabel();
if (tab != null && l.equalsIgnoreCase(tab.strip())) {
return e.getKey();
@@ -189,7 +190,7 @@ public final class LeadLauncher {
}
/** Start one lead. Returns false (having logged) rather than throwing on any failure. */
private boolean launch(String name, BridgedConfig.Leader lead, BridgedConfig.Profile profile) {
private boolean launch(String name, FleetConfig.Leader lead, FleetConfig.Profile profile) {
String label = lead.tabLabel();
String cwd = (lead.cwd() == null || lead.cwd().isBlank())
? System.getProperty("user.dir") : lead.cwd();
@@ -236,7 +237,7 @@ public final class LeadLauncher {
* The herdr agent kind for this profile — the same value the matching member adapter passes, so
* herdr resolves the same executable for a lead as it does for a member on that backend.
*/
private static String herdrKind(BridgedConfig.Profile profile) {
private static String herdrKind(FleetConfig.Profile profile) {
return profile.isOpenCode() ? "opencode" : "claude";
}
@@ -248,11 +249,12 @@ public final class LeadLauncher {
* any other primary. This is the single most important difference from the member launchers;
* do not "unify" it back.
*/
private List<String> leadArgv(BridgedConfig.Profile profile) {
private List<String> leadArgv(FleetConfig.Profile profile) {
List<String> argv = new ArrayList<>(profile.argv());
if (profile.hasMcp()) {
argv.add("--mcp-config");
argv.add("{\"mcpServers\":{\"bridge\":{\"type\":\"http\",\"url\":\""
argv.add("{\"mcpServers\":{\"" + PeerLauncher.MCP_MOUNT_NAME
+ "\":{\"type\":\"http\",\"url\":\""
+ profile.mcpUrl() + "\"}}}");
}
// Appended last, for the same reason the member launcher does it (CB-533): the argv is
@@ -274,7 +276,7 @@ public final class LeadLauncher {
* definition, so there is no configuration under which pointing it elsewhere is correct, and a
* profile that carries them (a member profile reused as a lead's backend) must not leak them in.
*/
private Map<String, String> leadEnv(BridgedConfig.Profile profile) {
private Map<String, String> leadEnv(FleetConfig.Profile profile) {
Map<String, String> out = new LinkedHashMap<>();
if (profile.env() != null) {
out.putAll(profile.env());
@@ -1,4 +1,4 @@
package dev.ltms.bridged.logging;
package dev.ltms.fleet.logging;
import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.Logger;
@@ -1,6 +1,6 @@
package dev.ltms.bridged.mcp;
package dev.ltms.fleet.mcp;
import dev.ltms.bridged.herdr.PaneLocator;
import dev.ltms.fleet.herdr.PaneLocator;
/**
* Resolves <em>who is calling</em> an MCP tool from the connection alone — the anti-spoofing
@@ -1,26 +1,28 @@
package dev.ltms.bridged.mcp;
package dev.ltms.fleet.mcp;
import dev.ltms.bridged.auth.AuditLog;
import dev.ltms.bridged.auth.Authz;
import dev.ltms.bridged.auth.CallerResolver;
import dev.ltms.bridged.auth.Principal;
import dev.ltms.bridged.auth.Role;
import dev.ltms.bridged.guard.GuardException;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.msg.Rendezvous;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.placement.BackendQuarantine;
import dev.ltms.bridged.placement.PlacementException;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.session.WorktreeRequest;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.peer.PeerLauncher;
import dev.ltms.fleet.auth.AuditLog;
import dev.ltms.fleet.auth.Authz;
import dev.ltms.fleet.auth.CallerResolver;
import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.auth.Role;
import dev.ltms.fleet.guard.GuardException;
import dev.ltms.fleet.herdr.Agent;
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;
import dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.WorktreeRequest;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.peer.PeerLauncher;
import io.modelcontextprotocol.common.McpTransportContext;
import io.modelcontextprotocol.json.McpJsonMapper;
import io.modelcontextprotocol.json.jackson3.JacksonMcpJsonMapperSupplier;
@@ -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;
@@ -61,9 +64,9 @@ import org.slf4j.LoggerFactory;
* {@link McpSchema.CallToolResult}, so it is unit-testable without standing up the HTTP transport;
* the SDK owns the wire protocol. Mount {@link #servlet()} at {@code /mcp} on the daemon's Jetty.
*/
public final class BridgeMcp {
public final class FleetMcp {
private static final Logger log = LoggerFactory.getLogger(BridgeMcp.class);
private static final Logger log = LoggerFactory.getLogger(FleetMcp.class);
private static final long DEFAULT_TIMEOUT_MS = 25_000;
private static final long MAX_TIMEOUT_MS = 120_000;
@@ -97,6 +100,8 @@ public final class BridgeMcp {
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<String, Integer> liveCount, Function<String, Integer> maxLoad,
@@ -127,10 +132,25 @@ public final class BridgeMcp {
* @param quarantine CB-578 stage B facts for {@code fleet_profiles}; required — pass
* {@link QuarantineSource#none()} for a caller that does not want the feature
*/
public BridgeMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
public FleetMcp(MessageService messages, PeerLauncher workers, SessionManager sessions,
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 BridgeMcp {
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 BridgeMcp {
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<McpSyncServerExchange, McpSchema.CallToolRequest, McpSchema.CallToolResult> stopHandler =
(exchange, req) -> {
@@ -377,7 +405,7 @@ public final class BridgeMcp {
* without fabricating an SDK {@code McpSyncServerExchange}.
*
* <p>This surface exists because the enforcement was previously unreachable from a test: no
* test constructs a {@code BridgeMcp}, so the whole MCP-side gate ran zero times in the suite
* test constructs a {@code FleetMcp}, so the whole MCP-side gate ran zero times in the suite
* while the REST-side equivalent had ten tests. A security control nothing exercises is a
* claim, not a control.
*
@@ -398,7 +426,7 @@ public final class BridgeMcp {
String reason = Authz.isUnauthenticated(caller) ? "unauthenticated" : "forbidden";
AuditLog.denied(caller, action, target, reason);
if (metrics != null) {
metrics.inc(BridgedMetrics.AUTH_FAILURES, "reason", reason);
metrics.inc(FleetMetrics.AUTH_FAILURES, "reason", reason);
}
return error(reason + ": " + caller.describe() + " may not " + action);
}
@@ -585,6 +613,54 @@ public final class BridgeMcp {
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.
*
* <p>{@code coordId} is mutually exclusive with {@code sessionId} and {@code turnId}: those two
* address a worker session owned by <em>this</em> 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.
*
* <p>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)) {
@@ -745,7 +821,7 @@ public final class BridgeMcp {
* CB-301-ext: {@code worktreeRequest} non-null provisions an isolated git worktree.
* CB-584: {@code sessionName}/{@code resumeSessionId} request agent session identity — a resumed
* conversation requires an explicit {@code profile} whose adapter declares
* {@link dev.ltms.bridged.peer.Capability#SESSION_RESUME}, or the spawn is refused rather than
* {@link dev.ltms.fleet.peer.Capability#SESSION_RESUME}, or the spawn is refused rather than
* silently starting a cold session.
*
* <p>{@code role} and {@code profile} are independent: the role picks the contract, the profile
@@ -870,6 +946,22 @@ public final class BridgeMcp {
static McpSchema.CallToolResult listFleet(PeerLauncher workers, SessionManager sessions, MessageService messages,
CapacitySource capacity, HealthCoverageSource healthCoverage,
QuarantineSource quarantine, Map<String, String> 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<String, String> leads, String selfTerm,
String selfCoordId) {
try {
Map<String, Agent> live = workers.list().stream()
.map(Agent.class::cast)
@@ -888,6 +980,9 @@ public final class BridgeMcp {
Map<String, Object> 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 BridgeMcp {
"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 BridgeMcp {
"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")));
}
@@ -1,4 +1,4 @@
package dev.ltms.bridged.mcp;
package dev.ltms.fleet.mcp;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.mcp;
package dev.ltms.fleet.mcp;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.mcp;
package dev.ltms.fleet.mcp;
/**
* Resolves the OS PID that owns a loopback TCP source port — the OS half of connection-based MCP
@@ -1,4 +1,4 @@
package dev.ltms.bridged.mcp;
package dev.ltms.fleet.mcp;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.mcp;
package dev.ltms.fleet.mcp;
/**
* Resolves a process's current working directory from its PID — the OS half of CB-112's
@@ -1,11 +1,12 @@
package dev.ltms.bridged.member;
package dev.ltms.fleet.member;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.guard.SubscriptionGuard;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.peer.Capability;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.PeerLauncher;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -53,7 +54,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
* existing deployments and tests keep the legacy non-blocking spawn semantics.
*/
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env) {
this(agents, spaces, guard, profiles, defaultProfile, env, 0,
System::currentTimeMillis, () -> sleepUninterruptibly(300));
@@ -64,7 +65,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
* until the pane reports an injectable state or {@code spawnReadyTimeoutMs} elapses.
*/
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs) {
this(agents, spaces, guard, profiles, defaultProfile, env,
@@ -77,10 +78,10 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
* <em>role</em>; a profile may still override it with its own {@code tabLabel}.
*/
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<BridgedConfig.Fleet> fleet) {
Supplier<FleetConfig.Fleet> fleet) {
this(agents, spaces, guard, profiles, defaultProfile, env,
spawnReadyTimeoutMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
@@ -91,11 +92,11 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
* Production constructor, plus the CB-596 {@code memberCredentials} policy supplier.
*/
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<BridgedConfig.Fleet> fleet,
Supplier<BridgedConfig.MemberCredentials> memberCredentials) {
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(agents, spaces, guard, profiles, defaultProfile, env,
spawnReadyTimeoutMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
@@ -120,7 +121,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
* (poll ms) is accepted for API symmetry but otherwise unused here
*/
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper) {
@@ -134,11 +135,11 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
* @param fleet live fleet config, read once for each spawn
*/
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<BridgedConfig.Fleet> fleet) {
Supplier<FleetConfig.Fleet> fleet) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet);
this.guard = guard;
@@ -148,12 +149,12 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
* Full testability constructor, plus the CB-596 {@code memberCredentials} policy supplier.
*/
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<BridgedConfig.Fleet> fleet,
Supplier<BridgedConfig.MemberCredentials> memberCredentials) {
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
this.guard = guard;
@@ -165,12 +166,12 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
* default (the real {@code System.getenv()} key set) via the constructor above.
*/
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<BridgedConfig.Fleet> fleet,
Supplier<BridgedConfig.MemberCredentials> memberCredentials,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<Set<String>> hostEnvNames) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames);
@@ -187,7 +188,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
* applied here — see {@link #applySessionIdentity}.
*/
@Override
protected Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec) {
protected Launch buildLaunch(FleetConfig.Profile cfg, LaunchSpec spec) {
// CB-539: a profile may deliberately opt into the subscription (subscription: true) when no
// off-subscription endpoint exists for it — e.g. `sonnet` on `ccs`. That profile gets no
// ANTHROPIC_BASE_URL/AUTH_TOKEN (there is nothing to point them at) and the guard's base_url
@@ -211,6 +212,17 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
guard.assertWorker(baseUrl); // hard stop before we spawn anything
}
// CB-634: the IDE guidance is delivered as an on-disk CLAUDE.local.md overlay, NOT through
// the charter — the charter returns to role -> reply only. Best-effort: a failed overlay
// must never fail the spawn, and `writeIdeOverlay` no-ops unless the cwd is a provisioned
// worktree (see its .git-file safety gate). The overlay pins, and the auto-open opens, the
// module dir (this repo's pom is in `bridged/`, not at the worktree root) — see ideProjectPath.
if (cfg.hasIdeMcp()) {
String projectPath = PeerLauncher.ideProjectPath(spec.cwd(), cfg.ideProjectDir());
writeIdeOverlay(spec.cwd(), projectPath);
PeerLauncher.openInIde(projectPath, cfg.ideOpenCommand(), log);
}
Map<String, String> workerEnv = baseEnv(cfg);
if (onSubscription) {
// CB-542 belt-and-braces: on the subscription path no guard vets these two keys, and the
@@ -232,11 +244,11 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
// handle is known BEFORE the agent has written anything; a resume spawn adopts its prior
// id via -r and passes no --session-id (the two conflict). Both are injected before the
// model flag so --model keeps outranking the operator's own argv.
// mutableArgv: argvWithBridge may hand back the profile's own (immutable) List.of when it
// mutableArgv: argvWithFleet may hand back the profile's own (immutable) List.of when it
// has neither MCP nor a charter — session flags must be added into a list we own.
List<String> argv = mutableArgv(argvWithBridge(cfg, spec));
List<String> 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);
}
/**
@@ -289,27 +301,33 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
* is present it keeps its proven inline {@code --append-system-prompt} delivery, which is also
* the only form that reaches a member with no repo checkout.
*/
private List<String> argvWithBridge(BridgedConfig.Profile cfg, LaunchSpec spec) {
private List<String> argvWithFleet(FleetConfig.Profile cfg, LaunchSpec spec) {
String roleCharter = nonBlank(spec.roleCharter());
// CB-634: the IDE guidance is delivered as an on-disk overlay (writeIdeOverlay), not through
// the charter. The charter file is role -> reply only.
String replyCharter = nonBlank(spec.replyCharter());
Path agentFile = agentDefinitionFile(spec.cwd(), spec.role(), ".claude", "agents");
if (!cfg.hasMcp() && roleCharter == null && replyCharter == null && agentFile == null) {
if (!cfg.mountsAnyMcp() && roleCharter == null
&& replyCharter == null && agentFile == null) {
return cfg.argv();
}
List<String> argv = mutableArgv(cfg.argv());
if (cfg.hasMcp()) {
String mcpJson = "{\"mcpServers\":{\"bridge\":{\"type\":\"http\",\"url\":\""
+ cfg.mcpUrl() + "\"}}}";
if (cfg.mountsAnyMcp()) {
argv.add("--mcp-config");
argv.add(mcpJson);
argv.add(mcpConfigJson(cfg));
}
if (roleCharter != null) {
String combined = replyCharter == null ? roleCharter : roleCharter + "\n\n" + replyCharter;
argv.add("--append-system-prompt-file");
argv.add(writeCharterFile(combined).toString());
} else if (replyCharter != null) {
// Combine the charters in order role -> reply, dropping any that are absent. When two
// or more survive they must ride one --append-system-prompt-file (CB-618 forbids the inline
// flag and the file flag together). A lone reply charter keeps its proven inline delivery.
List<String> charters = new java.util.ArrayList<>(2);
if (roleCharter != null) charters.add(roleCharter);
if (replyCharter != null) charters.add(replyCharter);
if (charters.size() == 1 && replyCharter != null && roleCharter == null) {
argv.add("--append-system-prompt");
argv.add(replyCharter);
} else if (!charters.isEmpty()) {
argv.add("--append-system-prompt-file");
argv.add(writeCharterFile(String.join("\n\n", charters)).toString());
}
if (agentFile != null) {
argv.add("--agent");
@@ -318,6 +336,83 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
return argv;
}
/**
* The {@code --mcp-config} JSON for this member: always the bridge mount when {@link
* FleetConfig.Profile#hasMcp()}, plus the IDE Index MCP as a second server named {@code
* intellij} when {@link FleetConfig.Profile#hasIdeMcp()} (CB-634). At least one is present —
* the caller only reaches here when {@link FleetConfig.Profile#mountsAnyMcp()} is true.
*/
private static String mcpConfigJson(FleetConfig.Profile cfg) {
StringBuilder servers = new StringBuilder();
if (cfg.hasMcp()) {
servers.append('"').append(PeerLauncher.MCP_MOUNT_NAME)
.append("\":{\"type\":\"http\",\"url\":\"").append(cfg.mcpUrl()).append("\"}");
}
if (cfg.hasIdeMcp()) {
if (servers.length() > 0) servers.append(',');
servers.append("\"intellij\":{\"type\":\"http\",\"url\":\"")
.append(cfg.ideMcpUrl()).append("\"}");
}
return "{\"mcpServers\":{" + servers + "}}";
}
/**
* Deliver the shared IDE guidance ({@link PeerLauncher#ideOverlayText}) as an on-disk
* {@code CLAUDE.local.md} overlay beside the project's own {@code CLAUDE.md} (CB-634), and
* register the overlay in the repository's common {@code info/exclude} so it never shows as
* untracked (git reads a worktree's excludes from the common dir, not the per-worktree gitdir).
*
* <p><strong>Safety gate:</strong> the overlay is written ONLY when {@code cwd/.git} is a
* <em>regular file</em> — a provisioned worktree keeps a {@code .git} FILE holding a
* {@code gitdir: <path>} pointer, while the primary's real checkout has a {@code .git}
* DIRECTORY. Returning without writing when {@code .git} is a directory is the whole safety of
* the feature: it must never write into a non-worktree cwd, i.e. never clobber a project that
* does not want the overlay.
*
* <p>Best-effort: a failure is logged at debug and swallowed — a failed overlay must never fail
* the spawn.
*
* @param cwd the member's worktree root, where the {@code CLAUDE.local.md} file is written
* @param projectPath the module dir the overlay pins {@code project_path} to (see
* {@link PeerLauncher#ideProjectPath}); equals {@code cwd} when no module subdir
*/
private static void writeIdeOverlay(String cwd, String projectPath) {
try {
Path dotGit = Path.of(cwd, ".git");
if (!Files.isRegularFile(dotGit)) {
// Not a provisioned worktree (primary's real checkout has a .git directory, or the
// cwd is not a repo at all). Never write into it.
return;
}
// The overlay FILE lives at the worktree root (claude-code's cwd), but its CONTENT pins
// project_path to the module dir the IDE opened (projectPath), not the worktree root.
Files.writeString(Path.of(cwd, "CLAUDE.local.md"), PeerLauncher.ideOverlayText(projectPath));
String gitdirLine = Files.readString(dotGit).trim();
Path gitDir = Path.of(gitdirLine.replaceFirst("^gitdir:\\s*", ""));
if (!gitDir.isAbsolute()) {
gitDir = Path.of(cwd).resolve(gitDir).normalize();
}
// git reads info/exclude from the COMMON dir, never the per-worktree gitdir (only
// info/sparse-checkout is per-worktree). A provisioned worktree's gitdir is
// <common>/worktrees/<name>, so the common dir is two levels up; writing the entry into
// the per-worktree gitdir leaves it un-honoured and the overlay shows as untracked.
Path commonDir = gitDir;
if (gitDir.getParent() != null && gitDir.getParent().getFileName() != null
&& "worktrees".equals(gitDir.getParent().getFileName().toString())) {
commonDir = gitDir.getParent().getParent();
}
Path exclude = commonDir.resolve("info").resolve("exclude");
Files.createDirectories(exclude.getParent());
String overlayLine = "CLAUDE.local.md";
if (!Files.exists(exclude) || Files.readAllLines(exclude).stream().noneMatch(overlayLine::equals)) {
Files.writeString(exclude, (Files.exists(exclude) ? System.lineSeparator() : "")
+ overlayLine + System.lineSeparator());
}
} catch (Exception e) {
log.debug("cannot write IDE overlay into worktree '{}'", cwd, e);
}
}
/** {@code s}, or {@code null} when {@code s} is null/blank — the charter-presence test used above. */
private static String nonBlank(String s) {
return (s == null || s.isBlank()) ? null : s;
@@ -358,7 +453,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
* {@code gx10} does) are untouched — this adds nothing when there is nothing to add. This is
* the {@code kind: claude} counterpart of the opencode adapter's {@code -m provider/model}.
*/
private static List<String> argvWithModel(List<String> argv, BridgedConfig.Profile cfg) {
private static List<String> argvWithModel(List<String> argv, FleetConfig.Profile cfg) {
if (cfg.model() == null || cfg.model().isBlank()) {
return argv;
}
@@ -368,6 +463,30 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
return withModel;
}
/**
* Pin a bounded auto-compaction window on the command line via {@code --autocompact <tokens>},
* 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).
*
* <p>Mirrors {@link #argvWithModel}: appended after it, so it survives the {@code ccs <profile>}
* 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 <auto|tokens>}
* (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<String> argvWithAutoCompact(List<String> argv, FleetConfig.Profile cfg) {
if (cfg.autoCompactWindow() == null) {
return argv;
}
List<String> 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. */
@@ -411,7 +530,7 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
/** Whether any configured profile opts into a git-forge token (required for {@link Capability#SELF_PR}). */
private boolean hasGitTokenProfile() {
return profileConfigs().stream().anyMatch(BridgedConfig.Profile::hasGitToken);
return profileConfigs().stream().anyMatch(FleetConfig.Profile::hasGitToken);
}
// --- CB-117 reap predicate (Claude prefix), kept for direct unit testing -------------------
@@ -1,19 +1,19 @@
package dev.ltms.bridged.member;
package dev.ltms.fleet.member;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.peer.Capability;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.peer.PeerHandle;
import dev.ltms.bridged.peer.PeerLauncher;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.peer.SpawnRequest;
import dev.ltms.bridged.placement.BackendQuarantine;
import dev.ltms.bridged.placement.PlacementCandidate;
import dev.ltms.bridged.placement.PlacementContext;
import dev.ltms.bridged.placement.PlacementException;
import dev.ltms.bridged.placement.PlacementPolicies;
import dev.ltms.bridged.placement.PlacementPolicy;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.Agent;
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 dev.ltms.fleet.placement.BackendQuarantine;
import dev.ltms.fleet.placement.PlacementCandidate;
import dev.ltms.fleet.placement.PlacementContext;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.placement.PlacementPolicies;
import dev.ltms.fleet.placement.PlacementPolicy;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -78,7 +78,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
* this way is the set of adapters ({@link #byProfile}), because a new backend needs a launcher
* and launchers are built once; {@code ConfigRef} classifies that as deferred and says so.
*/
private final Supplier<Map<String, BridgedConfig.Profile>> profileConfigs;
private final Supplier<Map<String, FleetConfig.Profile>> profileConfigs;
private final Supplier<PlacementPolicy> placementPolicy;
/** CB-578 stage B: credential cooldown, checked before an explicit spawn and filtered into placement. */
@@ -89,7 +89,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
* {@code null}, and an empty pool for a role, both fall back to every configured profile — the
* pre-CB-557 behaviour.
*/
private final Supplier<BridgedConfig.Fleet> fleet;
private final Supplier<FleetConfig.Fleet> fleet;
/**
* Backward-compatible constructor: fixed placement, no live-counting. Use this for tests and
@@ -119,7 +119,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
*/
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
String defaultProfile,
Map<String, BridgedConfig.Profile> profileConfigs,
Map<String, FleetConfig.Profile> profileConfigs,
PlacementPolicy placementPolicy,
Function<String, Integer> liveCount) {
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, null, BackendQuarantine.none());
@@ -136,27 +136,27 @@ public final class CompositePeerLauncher implements PeerLauncher {
*/
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
String defaultProfile,
Map<String, BridgedConfig.Profile> profileConfigs,
Map<String, FleetConfig.Profile> profileConfigs,
PlacementPolicy placementPolicy,
Function<String, Integer> liveCount,
BridgedConfig.Fleet fleet) {
FleetConfig.Fleet fleet) {
this(delegates, defaultProfile, profileConfigs, placementPolicy, liveCount, fleet, BackendQuarantine.none());
}
/**
* Production constructor with role pools and quarantine (CB-578 stage B). The full-featured
* non-reloading form; {@link #CompositePeerLauncher(List, String, Supplier, Function, BackendQuarantine)}
* is what {@code Bridged.main} actually wires up.
* is what {@code Fleetd.main} actually wires up.
*
* @param quarantine required — pass {@link BackendQuarantine#none()} for a caller that does not
* want the feature, never a defaulting overload (CB-578 stage B's own rule).
*/
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
String defaultProfile,
Map<String, BridgedConfig.Profile> profileConfigs,
Map<String, FleetConfig.Profile> profileConfigs,
PlacementPolicy placementPolicy,
Function<String, Integer> liveCount,
BridgedConfig.Fleet fleet,
FleetConfig.Fleet fleet,
BackendQuarantine quarantine) {
// LinkedHashMap, not Map.copyOf: candidates() promises definition order and the weighted
// policy breaks exact-weight ties on it, so a salted iteration order would make placement
@@ -175,7 +175,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
*/
public CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
String defaultProfile,
Supplier<BridgedConfig> config,
Supplier<FleetConfig> config,
Function<String, Integer> liveCount,
BackendQuarantine quarantine) {
this(delegates, defaultProfile,
@@ -189,10 +189,10 @@ public final class CompositePeerLauncher implements PeerLauncher {
/** The all-suppliers form every other constructor funnels into. */
private CompositePeerLauncher(List<HerdrPeerLauncher> delegates,
String defaultProfile,
Supplier<Map<String, BridgedConfig.Profile>> profileConfigs,
Supplier<Map<String, FleetConfig.Profile>> profileConfigs,
Supplier<PlacementPolicy> placementPolicy,
Function<String, Integer> liveCount,
Supplier<BridgedConfig.Fleet> fleet,
Supplier<FleetConfig.Fleet> fleet,
BackendQuarantine quarantine) {
this.fleet = fleet;
this.quarantine = Objects.requireNonNull(quarantine, "quarantine");
@@ -229,8 +229,8 @@ public final class CompositePeerLauncher implements PeerLauncher {
* <p>Read fresh on every call so a reload is visible; a caller that needs two consistent reads
* takes one local, as {@link #poolFor} does.
*/
private Map<String, BridgedConfig.Profile> profiles0() {
Map<String, BridgedConfig.Profile> m = profileConfigs.get();
private Map<String, FleetConfig.Profile> profiles0() {
Map<String, FleetConfig.Profile> m = profileConfigs.get();
return m == null ? Map.of() : m;
}
@@ -254,7 +254,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
String requestedProfile = req.profileName();
if (requestedProfile != null && !requestedProfile.isBlank()) {
// An explicit profile bypasses the placement policy, but not the capacity cap: maxLoad
// is documented as an unconditional limit on this profile (BridgedConfig.Profile), and
// is documented as an unconditional limit on this profile (FleetConfig.Profile), and
// the charter makes explicit-profile spawns the normal path — so skipping the check
// here would leave the cap dead config in real operation.
HerdrPeerLauncher d = route(requestedProfile);
@@ -317,7 +317,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
/**
* Refuse an explicit-profile spawn when the profile is at its {@code maxLoad} cap.
*
* <p>maxLoad is a documented, unconditional capacity limit (see {@code BridgedConfig.Profile#maxLoad}),
* <p>maxLoad is a documented, unconditional capacity limit (see {@code FleetConfig.Profile#maxLoad}),
* and the charter makes explicit-profile spawns the normal path — so enforcing it only in placement
* ({@code PlacementPolicyUtil}, package-private, hence not linked) would leave the cap dead config
* on every call that names a profile. Same rule as placement: {@code live >= cap} is at capacity.
@@ -359,7 +359,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
/** {@code profile}'s credential group (CB-578 stage B), or the profile's own name if unconfigured. */
private String credentialIdFor(String profile) {
BridgedConfig.Profile cfg = profiles0().get(profile);
FleetConfig.Profile cfg = profiles0().get(profile);
return cfg == null ? profile : cfg.effectiveCredentialId();
}
@@ -377,7 +377,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
// until CB-585: an explicit `maxLoad: 0` now survives as 0 and is a real cap of zero, so the
// check below refuses every spawn on that profile, and a negative value is refused at config
// load rather than normalized away.
BridgedConfig.Profile cfg = profiles0().get(profile);
FleetConfig.Profile cfg = profiles0().get(profile);
Integer cap = (cfg == null) ? null : cfg.maxLoad();
if (cap == null) {
return;
@@ -398,8 +398,8 @@ public final class CompositePeerLauncher implements PeerLauncher {
* a pool entry with no profile, so a survivor is a profile this particular composite does not own.
*/
private List<String> poolFor(MemberRole role) {
Map<String, BridgedConfig.Profile> configured = profiles0();
BridgedConfig.Fleet f = fleet.get();
Map<String, FleetConfig.Profile> configured = profiles0();
FleetConfig.Fleet f = fleet.get();
List<String> pool = (f == null) ? List.of() : f.profilesFor(role);
List<String> known = pool.stream().filter(configured::containsKey).toList();
return known.isEmpty() ? List.copyOf(configured.keySet()) : known;
@@ -415,7 +415,7 @@ public final class CompositePeerLauncher implements PeerLauncher {
private List<PlacementCandidate> candidates(MemberRole role) {
List<PlacementCandidate> out = new ArrayList<>();
for (String name : poolFor(role)) {
BridgedConfig.Profile w = profiles0().get(name);
FleetConfig.Profile w = profiles0().get(name);
if (w != null) {
out.add(new PlacementCandidate(name, null, w.weight(), w.maxLoad()));
}
@@ -0,0 +1,276 @@
package dev.ltms.fleet.member;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.IOException;
import java.io.UncheckedIOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.time.Duration;
import java.time.Instant;
import java.util.ArrayList;
import java.util.List;
import java.util.Set;
import java.util.stream.Stream;
/**
* CB-633: generates the per-spawn {@code ZDOTDIR} directory whose startup files enforce
* {@code memberCredentials.policy: allow-list}.
*
* <p>The seam: zsh reads its startup files from {@code $ZDOTDIR}, and the daemon puts that variable
* in the pane-creation env map. The operator's whole chain ({@code ~/.zshrc} → secret store) runs
* inside those files, so a scrub appended to the LAST one runs after everything the operator
* sourced, and nothing later can re-export over it. This is the property CB-596's env-overlay
* control lacked: herdr applies that overlay BEFORE the shell starts, so any sourced file can undo
* it — and did.
*
* <p><b>Which file is last depends on the platform, so the scrub runs from two of them.</b> zsh
* reads {@code .zshenv} always, {@code .zprofile} and {@code .zlogin} only for a LOGIN shell, and
* {@code .zshrc} only for an INTERACTIVE one. herdr does not open the same kind of shell
* everywhere — measured on herdr 0.8.0: macOS panes run {@code -zsh} (login, so {@code .zlogin}
* runs), Linux panes run a plain {@code /usr/bin/zsh} (interactive but NOT login, so
* {@code .zlogin} never runs at all). A scrub in {@code .zlogin} alone is therefore a control that
* silently does nothing on Linux — the exact failure this class exists to remove, one platform
* over.
*
* <p>So both {@code .zshrc} and {@code .zlogin} source the same generated {@code scrub.zsh} after
* sourcing their {@code $HOME} counterpart. On Linux only the first fires; on macOS both do, and
* the second pass is deliberate rather than merely harmless — it re-scrubs anything the operator's
* own {@code ~/.zlogin} exported after {@code .zshrc} had finished. Re-running is idempotent: a
* name already blank is blanked again, and the report is rewritten with the same counts.
*
* <p>Each generated file sources its {@code $HOME} counterpart FIRST, so {@code PATH} and every
* toolchain binary still resolve exactly as the operator configured them; only afterwards does
* {@code .zlogin} run the scrub: every EXPORTED variable not on the derived allow-list is re-exported
* blank. Blank, not credential-shaped-pattern-filtered: a pattern list ({@code *TOKEN*}, …) is an
* enumeration and misses what it did not think of — a username is the other half of a credential and
* is shaped like none. Credential-SHAPED names among the blanked set go to the WARN log only,
* never to the control.
*
* <p>The scrub also writes {@code scrub-report.txt} into its own directory: one {@code allowed N of
* M} line (N = exports left untouched, M = exports present when the scrub ran), then the blanked
* NAMES — never values. The launcher reads this back at teardown and logs it, because a blocked
* count next to an unknown denominator is not a finding.
*/
public final class EnvAllowListScrub {
private static final Logger log = LoggerFactory.getLogger(EnvAllowListScrub.class);
/** Name of the report file written into the generated directory by the scrub itself. */
static final String REPORT_FILE = "scrub-report.txt";
/** The scrub body, generated once and sourced from both {@code .zshrc} and {@code .zlogin}. */
static final String SCRUB_FILE = "scrub.zsh";
/** Prefix of every generated directory — also what {@link #reapOrphans} matches on. */
static final String DIR_PREFIX = "bridged-zdotdir-";
/**
* How old an orphan must be before {@link #reapOrphans} removes it. Comfortably longer than any
* spawn takes, so a directory belonging to a pane that is still starting is never removed.
*/
private static final Duration ORPHAN_AGE = Duration.ofHours(24);
/** Appended to the two startup files that must run the scrub, after their {@code $HOME} source. */
private static final String SOURCE_SCRUB =
"source \"$ZDOTDIR/" + SCRUB_FILE + "\"\n";
private EnvAllowListScrub() {
}
/**
* A parsed {@code scrub-report.txt}: how many exported variables existed when the scrub ran,
* how many were left untouched (allowed), and the NAMES that were blanked. Values never appear.
*/
record ScrubReport(int allowed, int total, List<String> blanked) {
}
/**
* Create a fresh ZDOTDIR directory under {@code parentDir} holding the four zsh startup files.
* Every file (and the directory) registers {@code deleteOnExit}, next to the existing per-spawn
* charter/config temp cleanup; the launcher additionally deletes eagerly at pane release.
*
* @param allowedNames the DERIVED allow-list — exact variable names that must survive the scrub
* @return the directory path (to be passed as the pane's {@code ZDOTDIR})
* @throws UncheckedIOException when the directory or any file cannot be written — a spawn whose
* protection cannot even be materialized must fail loudly rather
* than start unprotected
*/
public static Path generate(Path parentDir, Set<String> allowedNames) {
try {
reapOrphans(parentDir);
Path dir = Files.createTempDirectory(parentDir, DIR_PREFIX);
dir.toFile().deleteOnExit();
// The report is written by zsh, after these hooks are registered, so register its path
// too — otherwise the directory is non-empty at JVM exit and cannot be removed at all.
dir.resolve(REPORT_FILE).toFile().deleteOnExit();
write(dir, SCRUB_FILE, scrubScript(allowedNames));
write(dir, ".zshenv", homeSourcingFile(".zshenv"));
write(dir, ".zprofile", homeSourcingFile(".zprofile"));
write(dir, ".zshrc", homeSourcingFile(".zshrc") + SOURCE_SCRUB);
write(dir, ".zlogin", homeSourcingFile(".zlogin") + SOURCE_SCRUB);
return dir;
} catch (IOException e) {
throw new UncheckedIOException("cannot generate ZDOTDIR scrub files under " + parentDir, e);
}
}
/** One operator-sourcing startup file: source the {@code $HOME} counterpart, change nothing else. */
private static String homeSourcingFile(String name) {
return """
# generated by fleetd (CB-633 memberCredentials policy=allow-list) — do not edit.
# Source the operator's own %s first, so PATH and the agent binaries resolve as usual.
[ -r "$HOME/%s" ] && source "$HOME/%s"
""".formatted(name, name, name);
}
/**
* The generated {@code scrub.zsh} — the scrub body on its own, so the two startup files that
* must run it ({@code .zshrc} and {@code .zlogin}) hold one copy between them rather than two
* that can drift. Package-private so tests can assert on the exact script handed to zsh — the
* artefact here IS a shell file, and a test that checks only the Java string assembly proves
* nothing about whether zsh accepts it.
*/
static String scrubScript(Set<String> allowedNames) {
StringBuilder names = new StringBuilder();
for (String n : allowedNames.stream().sorted().toList()) {
if (names.length() > 0) {
names.append(' ');
}
// Names are validated against [A-Za-z_][A-Za-z0-9_]* before they get here; single quotes
// keep even a non-conforming name inert rather than executable.
names.append('\'').append(n.replace("'", "")).append('\'');
}
return """
# generated by fleetd (CB-633 memberCredentials policy=allow-list) — do not edit.
# Sourced from .zshrc and again from .zlogin, each time AFTER that file has sourced
# its $HOME counterpart — so this runs after everything the operator sourced, on a
# login shell (macOS panes) and on a plain interactive one (Linux panes) alike.
# Running twice is idempotent and deliberate: the second pass catches anything
# ~/.zlogin exported after ~/.zshrc had finished.
typeset -A _cb633_allowed
for _cb633_n in %s; do _cb633_allowed[$_cb633_n]=1; done
# Enumerate EXPORTED variable NAMES from `env` itself. Deliberately NOT the special
# `parameters` assoc: its subscript is evaluated arithmetically on this host's zsh
# and blows up on some names ("bad math expression"). Names not matching the
# identifier pattern (junk from multi-line values) are skipped, never scrubbed.
typeset -a _cb633_names
_cb633_names=("${(@f)$(command env | command cut -d= -f1)}")
typeset -a _cb633_blank
_cb633_blank=()
integer _cb633_total=0
for _cb633_n in "${_cb633_names[@]}"; do
[[ "$_cb633_n" =~ ^[A-Za-z_][A-Za-z0-9_]*$ ]] || continue
(( _cb633_total += 1 ))
[[ -n "${_cb633_allowed[$_cb633_n]-}" ]] && continue
case "$_cb633_n" in %s) continue ;; esac
_cb633_blank+=("$_cb633_n")
done
{ for _cb633_n in "${_cb633_blank[@]}"; do export "$_cb633_n="; done; } 2>/dev/null
integer _cb633_kept=$(( _cb633_total - ${#_cb633_blank} ))
{
print -r -- "allowed $_cb633_kept of $_cb633_total"
for _cb633_n in "${_cb633_blank[@]}"; do print -r -- "$_cb633_n"; done
} > "$ZDOTDIR/%s" 2>/dev/null
unset _cb633_allowed _cb633_names _cb633_blank _cb633_n _cb633_total _cb633_kept
""".formatted(names, MemberEnvAllowList.zshCasePattern(), REPORT_FILE);
}
private static void write(Path dir, String fileName, String content) throws IOException {
Path file = dir.resolve(fileName);
Files.writeString(file, content);
file.toFile().deleteOnExit();
}
/**
* Read and parse {@link #REPORT_FILE} out of a generated ZDOTDIR directory. Returns {@code null}
* when absent or unreadable (the pane may have been torn down before its login shell ever got to
* the scrub) — callers treat that as "no measurement available", never as success.
*/
static ScrubReport readReport(Path zdotdir) {
Path report = zdotdir.resolve(REPORT_FILE);
if (!Files.isRegularFile(report)) {
return null;
}
try {
List<String> lines = Files.readAllLines(report);
if (lines.isEmpty() || !lines.getFirst().startsWith("allowed ")) {
return null;
}
String[] parts = lines.getFirst().substring("allowed ".length()).trim().split("\\s+");
if (parts.length != 3 || !"of".equals(parts[1])) {
return null;
}
List<String> blanked = new ArrayList<>();
for (int i = 1; i < lines.size(); i++) {
if (!lines.get(i).isBlank()) {
blanked.add(lines.get(i));
}
}
return new ScrubReport(Integer.parseInt(parts[0]), Integer.parseInt(parts[2]),
List.copyOf(blanked));
} catch (IOException | NumberFormatException e) {
return null;
}
}
/** Best-effort recursive delete; failures are swallowed — JVM-exit cleanup is the backstop. */
/**
* Remove generated directories left behind by an earlier daemon process.
*
* <p>{@link #generate} registers each directory for deletion at JVM exit, which covers a clean
* shutdown and covers nothing else. A {@code kill -9}, a crash, or a host reboot leaves the
* directory in the temp dir for good, and the daemon is restarted often enough that these
* accumulate. They hold no secrets — the generated files contain variable NAMES and a report of
* names, never a value — but an unbounded pile of them in {@code /tmp} is still our mess to
* clear.
*
* <p>Called from {@link #generate}, so it runs on the path that creates them and needs no
* separate wiring or scheduler. Only directories older than {@link #ORPHAN_AGE} are touched,
* which keeps it clear of any pane that is merely still starting, including one belonging to a
* different daemon instance running right now. Best-effort: every failure is ignored, because
* tidying temp files must never be the reason a spawn fails.
*/
static void reapOrphans(Path parentDir) {
Instant cutoff = Instant.now().minus(ORPHAN_AGE);
try (Stream<Path> entries = Files.list(parentDir)) {
entries.filter(d -> d.getFileName().toString().startsWith(DIR_PREFIX))
.filter(Files::isDirectory)
.filter(d -> olderThan(d, cutoff))
.forEach(EnvAllowListScrub::deleteRecursively);
} catch (IOException | RuntimeException e) {
log.debug("could not scan {} for orphaned ZDOTDIRs: {}", parentDir, e.toString());
}
}
private static boolean olderThan(Path dir, Instant cutoff) {
try {
return Files.getLastModifiedTime(dir).toInstant().isBefore(cutoff);
} catch (IOException e) {
return false; // unreadable timestamp ⇒ leave it alone
}
}
static void deleteRecursively(Path dir) {
if (dir == null || !Files.exists(dir)) {
return;
}
try (Stream<Path> walk = Files.walk(dir)) {
walk.sorted(java.util.Comparator.reverseOrder()).forEach(p -> {
try {
Files.deleteIfExists(p);
} catch (IOException ignored) {
// best effort — deleteOnExit retries at JVM shutdown
}
});
} catch (IOException ignored) {
// same
}
}
}
@@ -1,19 +1,19 @@
package dev.ltms.bridged.member;
package dev.ltms.fleet.member;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.herdr.Tab;
import dev.ltms.bridged.herdr.Workspace;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.peer.Capability;
import dev.ltms.bridged.peer.CharterReceipt;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.peer.PeerHandle;
import dev.ltms.bridged.peer.PeerLauncher;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.peer.SpawnRequest;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.herdr.Tab;
import dev.ltms.fleet.herdr.Workspace;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.CharterReceipt;
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.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -49,7 +49,7 @@ import java.util.regex.Pattern;
* <li>{@code namePrefix} (constructor arg) — the label prefix ({@code claude}, {@code opencode})
* that drives both unique naming and the orphan-reap pattern, so each adapter reaps only its
* own kind of pane and never another's.</li>
* <li>{@link #buildLaunch(BridgedConfig.Profile, LaunchSpec)} — the peer-specific env map + argv, including any
* <li>{@link #buildLaunch(FleetConfig.Profile, LaunchSpec)} — the peer-specific env map + argv, including any
* subscription/guard check, MCP mount, and instruction injection. The base never sees how the
* peer is configured; it only places and starts the returned {@link Launch}.</li>
* </ul>
@@ -76,7 +76,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
private final String namePrefix; // label prefix: naming + reap scheme
private final AgentControl agents;
private final WorkspaceControl spaces;
private final Map<String, BridgedConfig.Profile> profiles; // profile name → spawn settings
private final Map<String, FleetConfig.Profile> profiles; // profile name → spawn settings
private final String defaultProfile; // profile a no-arg spawn uses (nullable)
/** Host env lookup (injectable for tests); adapters read it in {@link #buildLaunch}. */
@@ -91,14 +91,14 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* <p>CB-559: a supplier rather than a snapshot, so a config reload affects the next launch
* without a restart. Existing tabs keep the label they were given.
*/
private final Supplier<BridgedConfig.Fleet> fleet;
private final Supplier<FleetConfig.Fleet> fleet;
/**
* CB-596: the live {@code memberCredentials:} policy, read once per spawn (same hot-reload shape
* as {@link #fleet}). {@code null} — either the supplier itself, or what it returns — means no
* policy is configured and {@link #applyMemberCredentialPolicy} shadows nothing.
*/
private final Supplier<BridgedConfig.MemberCredentials> memberCredentials;
private final Supplier<FleetConfig.MemberCredentials> memberCredentials;
/**
* Enumerates the daemon's own process environment variable NAMES ONLY, never values — the CB-596
@@ -156,6 +156,20 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
private final ConcurrentMap<String, String> paneByAgentId = new ConcurrentHashMap<>();
private final AtomicBoolean resetUnsupportedLogged = new AtomicBoolean();
/**
* CB-633: the per-spawn ZDOTDIR directory generated for a pane under
* {@code memberCredentials.policy: allow-list}, keyed by herdr pane id so every teardown exit
* ({@link #stop} is reached from explicit DELETE, orphan reap, and the spawn-readiness gate
* timeout alike) can read the scrub's own report and then remove the directory. A pane that
* never reaches {@code stop} (spawn failure) leaks its directory only until JVM exit, where the
* generator's {@code deleteOnExit} hooks are the backstop — the same cleanup shape the existing
* charter/config temp files use.
*/
private final ConcurrentMap<String, Path> zdotdirByPane = new ConcurrentHashMap<>();
/** Guards {@link #warnNonZsh} to one WARN per launcher instance, not one per spawn. */
private final AtomicBoolean nonZshShellWarned = new AtomicBoolean();
/**
* @param namePrefix label prefix for this peer kind (drives naming and reap)
* @param agents herdr agent control (start, status, close)
@@ -169,7 +183,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* interval is baked into this hook, so the base needs no poll field
*/
protected HerdrPeerLauncher(String namePrefix, AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper) {
@@ -185,11 +199,11 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* above, so every existing call site keeps the default without an edit.
*/
protected HerdrPeerLauncher(String namePrefix, AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<BridgedConfig.Fleet> fleet) {
Supplier<FleetConfig.Fleet> fleet) {
this(namePrefix, agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
nowMillis, sleeper, fleet, null);
}
@@ -203,12 +217,12 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* existing call site keeps the pre-CB-596 default without an edit.
*/
protected HerdrPeerLauncher(String namePrefix, AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<BridgedConfig.Fleet> fleet,
Supplier<BridgedConfig.MemberCredentials> memberCredentials) {
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(namePrefix, agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
nowMillis, sleeper, fleet, memberCredentials, null);
}
@@ -218,12 +232,12 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* seam only — every production call site leaves this {@code null} and gets the real host env.
*/
protected HerdrPeerLauncher(String namePrefix, AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Supplier<BridgedConfig.Fleet> fleet,
Supplier<BridgedConfig.MemberCredentials> memberCredentials,
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials,
Supplier<Set<String>> hostEnvNames) {
this.fleet = fleet;
this.namePrefix = namePrefix;
@@ -246,7 +260,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* Any subscription/guard check, MCP mount, and instruction injection happen here. The env map
* and argv are adapter-private; the base only places and starts what is returned.
*/
protected abstract Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec);
protected abstract Launch buildLaunch(FleetConfig.Profile cfg, LaunchSpec spec);
/** Direct transport access for peer-specific, non-turn control operations. */
protected final AgentControl agents() {
@@ -327,7 +341,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
if (name == null || name.isBlank()) {
return List.of();
}
BridgedConfig.Profile cfg = profiles.get(name);
FleetConfig.Profile cfg = profiles.get(name);
return cfg == null ? List.of() : cfg.parityOverlay();
}
@@ -351,18 +365,18 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
/** The configured profiles, for adapter capability decisions (e.g. any git-token grant). */
protected Collection<BridgedConfig.Profile> profileConfigs() {
protected Collection<FleetConfig.Profile> profileConfigs() {
return profiles.values();
}
/** Resolve {@code profileName} (null/blank → default) to its config, or throw with the options. */
protected BridgedConfig.Profile requireProfile(String profileName) {
protected FleetConfig.Profile requireProfile(String profileName) {
String name = (profileName == null || profileName.isBlank()) ? defaultProfile : profileName;
if (name == null || name.isBlank()) {
throw new IllegalArgumentException("no default worker profile is configured — "
+ "pass a profile; configured: " + profiles.keySet());
}
BridgedConfig.Profile cfg = profiles.get(name);
FleetConfig.Profile cfg = profiles.get(name);
if (cfg == null) {
throw new IllegalArgumentException("unknown worker profile '" + name
+ "' — configured: " + profiles.keySet());
@@ -399,14 +413,14 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
/**
* Spawn a peer with session identity (CB-547a). The session values, role, and charter are
* threaded from the {@link SpawnRequest} into {@link #buildLaunch(BridgedConfig.Profile,
* threaded from the {@link SpawnRequest} into {@link #buildLaunch(FleetConfig.Profile,
* LaunchSpec)}, and the launch's resolved agent-session id is returned alongside the agent so
* the caller can put it on the {@link PeerHandle}.
*/
protected Spawned spawnInternal(String profileName, String requestedCwd, String callerCwd,
String sessionName, String resumeSessionId, MemberRole role) {
BridgedConfig.Profile cfg = requireProfile(profileName);
BridgedConfig.Fleet liveFleet = fleet == null ? null : fleet.get();
FleetConfig.Profile cfg = requireProfile(profileName);
FleetConfig.Fleet liveFleet = fleet == null ? null : fleet.get();
String roleCharter = liveFleet == null ? null : liveFleet.charterFor(role);
String replyCharter = cfg.hasMcp() ? REPLY_CHARTER : null;
String charter = roleCharter == null ? replyCharter
@@ -420,9 +434,27 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
try {
Launch launch = buildLaunch(cfg, new LaunchSpec(sessionName, resumeSessionId, role, charter,
roleCharter, replyCharter, cwd));
Agent agent = cfg.tabPlacement()
? spawnInTab(cfg, launch.env(), launch.argv(), cwd, role, liveFleet)
: spawnAsPane(cfg, launch.env(), launch.argv(), cwd, charter);
// CB-633: applied AFTER buildLaunch so the generated scrub's allow-list can also cover
// the exact env-map keys this launch injects (ANTHROPIC_*, OPENCODE_CONFIG, GITEA_TOKEN, …)
// — anything the daemon deliberately sets must survive its own control.
Path zdotdir = applyEnvironmentAllowListPolicy(cfg, launch);
Agent agent;
try {
agent = cfg.tabPlacement()
? spawnInTab(cfg, launch.env(), launch.argv(), cwd, role, liveFleet)
: spawnAsPane(cfg, launch.env(), launch.argv(), cwd, charter);
} catch (RuntimeException e) {
// A failed spawn has no pane id, so nothing would ever key this directory for
// teardown and it would sit in the temp dir until the JVM exits cleanly — which,
// for a daemon, may be never. Remove it on the way out.
if (zdotdir != null) {
EnvAllowListScrub.deleteRecursively(zdotdir);
}
throw e;
}
if (zdotdir != null) {
zdotdirByPane.put(agent.paneId(), zdotdir);
}
logCharterReceipt(receipt, true);
return new Spawned(agent, launch.agentSessionId(), receipt);
} catch (RuntimeException e) {
@@ -507,7 +539,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* guaranteed last resort so a pathological environment with an unset {@code user.dir} still
* honours the "never assume {@code $HOME}" contract rather than letting herdr default the pane.
*/
private static String resolveCwd(String requestedCwd, BridgedConfig.Profile cfg, String callerCwd) {
private static String resolveCwd(String requestedCwd, FleetConfig.Profile cfg, String callerCwd) {
return firstNonBlank(requestedCwd, cfg.cwd(), callerCwd, System.getProperty("user.dir"), ".");
}
@@ -519,8 +551,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
/** Dedicated worker space → own tab (carrying cwd+env) → start the peer into the seed pane. */
private Agent spawnInTab(BridgedConfig.Profile cfg, Map<String, String> workerEnv,
List<String> argv, String cwd, MemberRole role, BridgedConfig.Fleet liveFleet) {
private Agent spawnInTab(FleetConfig.Profile cfg, Map<String, String> workerEnv,
List<String> argv, String cwd, MemberRole role, FleetConfig.Fleet liveFleet) {
Workspace space = spaces.ensureWorkspace(cfg.workspace());
Tab.Created tab = spaces.createTab(space.workspaceId(), cwd, workerEnv);
log.info("spawning {} profile={} space={} tab={} cwd={}",
@@ -578,7 +610,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* composed charter, if any; its argv element is replaced by its digest so the log still shows
* which args were passed without exposing the charter prose.
*/
private Agent spawnAsPane(BridgedConfig.Profile cfg, Map<String, String> workerEnv,
private Agent spawnAsPane(FleetConfig.Profile cfg, Map<String, String> workerEnv,
List<String> argv, String cwd, String charter) {
log.info("spawning {} (pane placement) profile={} cwd={} argv={}",
namePrefix, cfg.profile(), cwd, redactCharter(argv, charter));
@@ -620,7 +652,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* backstop for the astronomically unlikely nonce+seq clash; the name is a label only — herdr
* detects kind and status from terminal output, not from it.
*/
private Started startUniquelyNamed(BridgedConfig.Profile cfg, List<String> argv, String paneId) {
private Started startUniquelyNamed(FleetConfig.Profile cfg, List<String> argv, String paneId) {
// Protocol 19 resolves the executable from the agent kind (== namePrefix here), so
// argv[0] — the configured executable — is dropped and only the extra args are passed.
List<String> args = argv.isEmpty() ? argv : argv.subList(1, argv.size());
@@ -774,11 +806,19 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
log.debug("not closing tab {} — it holds {} panes (not a dedicated peer tab)",
loc.tabId(), loc.tabPaneCount());
}
// CB-633: log this pane's allowed-N-of-M scrub report, then remove the generated ZDOTDIR.
// Last in, best-effort — a failure here must not mask a real teardown failure above.
try {
releaseZdotdir(paneId);
} catch (RuntimeException e) {
log.warn("memberCredentials allow-list: releasing ZDOTDIR for pane {} failed: {}",
paneId, e.getMessage());
}
}
/** Whether any configured profile places peers in their own tab (so tabs may need cleanup). */
private boolean usesTabPlacement() {
return profiles.values().stream().anyMatch(BridgedConfig.Profile::tabPlacement);
return profiles.values().stream().anyMatch(FleetConfig.Profile::tabPlacement);
}
/** True when a herdr error means the target is already gone (safe to treat as done). */
@@ -844,7 +884,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* {@code GITEA_HOST}. Push over SSH is unaffected; the only incremental grant is PR-create.
* Peer-neutral, so every herdr adapter reuses it unchanged.
*/
protected void applyGitToken(Map<String, String> workerEnv, BridgedConfig.Profile cfg) {
protected void applyGitToken(Map<String, String> workerEnv, FleetConfig.Profile cfg) {
if (!cfg.hasGitToken()) {
return;
}
@@ -860,7 +900,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* herdr's own login-shell process environment (gitea issue #82, superseding CB-592's single
* hardcoded {@code GITEA_ACCESS_TOKEN} name — see {@link #applyMemberCredentialPolicy}). herdr
* spawns a pane from its <em>own</em> process environment and layers our map on top —
* {@link dev.ltms.bridged.herdr.WorkspaceControl#createTab} and {@code #splitPane} send only
* {@link dev.ltms.fleet.herdr.WorkspaceControl#createTab} and {@code #splitPane} send only
* the keys we put in that map, so any key we never mention passes straight through from
* herdr's own shell, admin credentials included.
*
@@ -922,7 +962,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* <p>Adapter-specific variables are layered on top of this by {@code buildLaunch} and therefore
* win. That ordering is deliberate and load-bearing: it stops a profile's {@code env:} from
* overriding {@code ANTHROPIC_BASE_URL} and slipping past {@link
* dev.ltms.bridged.guard.SubscriptionGuard}, which is checked against the profile's
* dev.ltms.fleet.guard.SubscriptionGuard}, which is checked against the profile's
* {@code baseUrl} and nothing else.
*
* <p>The CB-596 credential shadow and the CB-592 marker are put in <em>last</em>, after the
@@ -931,7 +971,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* place both are applied: every {@code buildLaunch} in every adapter calls this first, so a new
* profile, and a peer kind not yet written, gets them for free.
*/
protected Map<String, String> baseEnv(BridgedConfig.Profile cfg) {
protected Map<String, String> baseEnv(FleetConfig.Profile cfg) {
Map<String, String> workerEnv = new LinkedHashMap<>();
String path = env.apply("PATH");
if (path != null && !path.isBlank()) {
@@ -955,18 +995,136 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
*
* <p>No {@code memberCredentials} configured — the supplier is {@code null}, or it resolves to
* one whose {@code known} list is empty — shadows nothing. This is a real, config-driven gap
* (see {@link BridgedConfig.MemberCredentials}'s javadoc), not a safe default: deny-by-default
* (see {@link FleetConfig.MemberCredentials}'s javadoc), not a safe default: deny-by-default
* only defends names the operator has actually enumerated in {@code known}.
*
* <p>CB-633: under {@code policy: allow-list} this overlay is NOT the control anymore — it is
* applied before the login shell runs and a sourced file can (and did) undo it. The control is
* the ZDOTDIR scrub ({@link #applyEnvironmentAllowListPolicy}); {@code known}/{@code allow}
* remain as reporting only via {@link #logCredentialGap}.
*/
private void applyMemberCredentialPolicy(Map<String, String> workerEnv) {
BridgedConfig.MemberCredentials creds = memberCredentials == null ? null : memberCredentials.get();
FleetConfig.MemberCredentials creds = memberCredentials == null ? null : memberCredentials.get();
if (creds == null) {
return;
}
if (!creds.isAllowList()) {
overlayBlockedCredentials(workerEnv, creds);
}
logCredentialGap(creds);
}
/** Put {@link #BLOCKED_CREDENTIAL_SENTINEL} over every blocked name in the pane-creation env map. */
private static void overlayBlockedCredentials(Map<String, String> workerEnv,
FleetConfig.MemberCredentials creds) {
for (String name : creds.blockedSet()) {
workerEnv.put(name, BLOCKED_CREDENTIAL_SENTINEL);
}
logCredentialGap(creds);
}
/**
* CB-633: under {@code memberCredentials.policy: allow-list}, generate the per-spawn ZDOTDIR
* directory whose startup files blank every exported variable not on the DERIVED allow-list —
* running AFTER the pane's shell has finished sourcing the operator's chain, which is what no
* pre-shell env overlay can achieve. {@code EnvAllowListScrub} sources the scrub from both the
* generated {@code .zshrc} and {@code .zlogin}, since a herdr pane is a login shell on macOS
* and a plain interactive one on Linux. Mutates {@code launch.env()} to carry
* {@code ZDOTDIR=<dir>}, so both placement paths ({@link #spawnInTab}, {@link #spawnAsPane})
* pass it through {@code tab.create}/{@code pane.split}. Returns the directory for teardown
* registration, or {@code null} when the policy does not apply.
*
* <p>The allow-list handed to the generator is the derived profile set ({@link
* MemberEnvAllowList#derive}) UNIONed with the exact keys of THIS launch's env map — names the
* daemon itself injects must survive its own control. {@code SSH_AUTH_SOCK} is added ONLY when
* the config explicitly allows it; by default it is absent, so the scrub blanks it like any
* other non-derived name.
*/
private Path applyEnvironmentAllowListPolicy(FleetConfig.Profile cfg, Launch launch) {
FleetConfig.MemberCredentials creds = memberCredentials == null ? null : memberCredentials.get();
if (creds == null || !creds.isAllowList()) {
return null;
}
String loginShell = resolveEnv("SHELL");
boolean zsh = loginShell != null && (loginShell.endsWith("/zsh") || loginShell.equals("zsh"));
if (!zsh) {
// A non-zsh login shell ignores ZDOTDIR entirely: NO scrub would run, so pretending
// otherwise would be worse than saying so. Warn loudly and fall back to the CB-596
// sentinel overlay over the enumerated known: names — weaker (a sourced file can undo
// it), but strictly better than nothing.
warnNonZsh(loginShell);
overlayBlockedCredentials(launch.env(), creds);
logCredentialGap(creds);
return null;
}
Set<String> allowed = new java.util.TreeSet<>(MemberEnvAllowList.derive(profiles.values()));
if (creds.sshAuthSockAllowed()) {
allowed.add(SSH_AUTH_SOCK);
} // blocked by default: absent from the set ⇒ blanked by the scrub like any other name
allowed.addAll(launch.env().keySet());
Path dir = EnvAllowListScrub.generate(Path.of(System.getProperty("java.io.tmpdir")), allowed);
launch.env().put("ZDOTDIR", dir.toAbsolutePath().toString());
log.info("memberCredentials policy=allow-list: profile={} generated ZDOTDIR {} — derived "
+ "allow-list holds {} name(s); the pane reports allowed N of M at release",
cfg.profile(), dir.getFileName(), allowed.size());
return dir;
}
/** The operator ssh-agent handle — kept ONLY by explicit config decision, never by default. */
private static final String SSH_AUTH_SOCK = "SSH_AUTH_SOCK";
/**
* CB-633: a non-zsh login shell means the allow-list control CANNOT run — say so once per
* launcher instance, naming the shell, instead of failing silently.
*/
private void warnNonZsh(String shell) {
if (nonZshShellWarned.compareAndSet(false, true)) {
log.warn("memberCredentials policy=allow-list: member login shell '{}' is NOT zsh — "
+ "ZDOTDIR scrubbing cannot run, so members' inherited environment is "
+ "UNPROTECTED beyond the enumerated known: fallback. Move herdr onto a "
+ "zsh account or switch policy back to deny-by-default.",
shell == null ? "<unset>" : shell);
}
}
/**
* CB-633 teardown half: read the pane's scrub report (the denominator report the generated
* scrub wrote) and delete the directory. Called from {@link #stop}, which is the one funnel
* every teardown exit already goes through.
*
* <p><b>A missing report is a WARN, not a debug line.</b> The report is the only evidence that
* the scrub ran at all in that pane. Its absence has an innocent reading — the pane died before
* its shell finished starting — and a serious one: the shell was not zsh, or it read its
* startup files from somewhere other than the directory we generated, in which case the member
* ran for its whole life with the operator's full secret store in its environment and nothing
* said so. We cannot tell those two apart from here, so the line says what is and is not known
* rather than picking one. Logging this at debug is how a control that silently stopped working
* stays unnoticed — the failure mode this whole class exists to remove.
*/
private void releaseZdotdir(String paneId) {
Path dir = zdotdirByPane.remove(paneId);
if (dir == null) {
return;
}
EnvAllowListScrub.ScrubReport report = EnvAllowListScrub.readReport(dir);
if (report == null) {
log.warn("memberCredentials allow-list: pane {} left no scrub report in {} — the "
+ "environment scrub cannot be confirmed to have run. Either the pane ended "
+ "before its shell finished starting, or its shell never read our generated "
+ "startup files, in which case that member saw the full host environment.",
paneId, dir);
} else {
log.info("memberCredentials allow-list: pane {} allowed {} of {} environment variables",
paneId, report.allowed(), report.total());
List<String> shaped = report.blanked().stream()
.filter(name -> CREDENTIAL_SHAPED_NAME.matcher(name).matches())
.toList();
if (!shaped.isEmpty()) {
log.warn("memberCredentials allow-list: pane {} blanked credential-shaped variable(s) "
+ "{} — confirm none of them was something a member legitimately needed",
paneId, shaped);
}
}
EnvAllowListScrub.deleteRecursively(dir);
}
/** Credential-shaped env var name heuristic for {@link #logCredentialGap} — case-insensitive. */
@@ -984,7 +1142,7 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* every such NAME, at WARN, at most once per launcher instance — never a value, a prefix of a
* value, or a hash of a value, so the log itself cannot leak anything.
*/
private void logCredentialGap(BridgedConfig.MemberCredentials creds) {
private void logCredentialGap(FleetConfig.MemberCredentials creds) {
Set<String> covered = new HashSet<>(creds.known());
covered.addAll(creds.allow());
List<String> gap = hostEnvNames.get().stream()
@@ -0,0 +1,112 @@
package dev.ltms.fleet.member;
import dev.ltms.fleet.config.FleetConfig;
import java.util.Collection;
import java.util.Set;
import java.util.TreeSet;
/**
* CB-633: the set of environment variable NAMES a spawned member is allowed to keep under
* {@code memberCredentials.policy: allow-list} — DERIVED from what the launcher itself injects,
* never hand-typed.
*
* <p>A hand-typed allow-list is the defect this class exists to prevent: a name an operator forgets
* to type is a credential that passes through to every member, and a profile added to config later
* would silently break spawns whose scrub did not know its names. Derivation closes both ends. The
* kept-name set is the union of:
*
* <ul>
* <li>every configured {@link FleetConfig.Profile profile}'s {@code gitTokenEnv},
* {@code gitHostEnv}, and {@code tokenEnv} values — these are variable <em>names</em> held in
* config, and the launcher reads their values out of exactly these variables;</li>
* <li>every key of every profile's {@code env:} map — anything the operator routes into a pane on
* purpose;</li>
* <li>{@link #INFRASTRUCTURE_PASSTHROUGH} — names that are not credentials at all but that a shell
* or the agent binary genuinely needs to function.</li>
* </ul>
*
* <p>Because the union spans EVERY profile (not just the one spawning), adding a new profile can
* only ever widen the list — it cannot break another spawn's scrub. And because the launcher also
* unions in the exact keys of each spawn's own env map at generation time (see {@code
* HerdrPeerLauncher}), anything the daemon deliberately injects for THIS spawn survives its own
* control.
*
* <p>{@code SSH_AUTH_SOCK} is deliberately NOT here. It is a handle to the operator's ssh-agent — a
* member holding it can sign with the operator's keys — so keeping it is a config decision
* ({@code memberCredentials.sshAuthSock: allow}), not a derivation default.
*/
public final class MemberEnvAllowList {
/**
* Names that are not credentials and that a login shell or agent binary genuinely needs.
*
* <p>Every name here is a location or a shell setting, never a credential. That rule is load
* bearing, and CB-633's first cut broke it: it also listed {@code ANTHROPIC_AUTH_TOKEN},
* {@code GITEA_TOKEN}, {@code GITEA_HOST}, {@code ANTHROPIC_BASE_URL}, {@code ANTHROPIC_MODEL},
* {@code CLAUDE_CONFIG_DIR}, {@code OPENCODE_CONFIG} and {@code BRIDGED_MEMBER} "because the
* launcher injects them". The launcher does — but only on the spawns where it actually sets
* them, and {@code HerdrPeerLauncher} already unions THIS spawn's env-map keys into the
* allow-list. So a static entry adds nothing on a spawn that injects the name, and on a spawn
* that does not it lets the operator's own value through under exactly the name a member reads.
* {@code ANTHROPIC_BASE_URL} is the sharpest case: an inherited one silently moves a member off
* the endpoint the profile chose.
*
* <p>{@code ZDOTDIR} stays because it is this control's own handle — lose it and every later
* sub-shell loses the scrub. {@code JAVA_HOME} and the {@code XDG_*} roots are toolchain
* locations. Everything else a member needs must arrive via a profile's {@code env:} or the
* launcher's own injection, both of which land on the derived set automatically.
*/
public static final Set<String> INFRASTRUCTURE_PASSTHROUGH = Set.of(
"PATH", "HOME", "SHELL", "TERM", "LANG", "TMPDIR",
"USER", "LOGNAME", "PWD", "SHLVL", "EDITOR", "PAGER",
"_",
"ZDOTDIR",
"JAVA_HOME",
"XDG_CONFIG_HOME", "XDG_DATA_HOME", "XDG_CACHE_HOME", "XDG_STATE_HOME");
/** Locale-category prefix kept as infrastructure ({@code LC_ALL}, {@code LC_CTYPE}, …). */
private static final String INFRASTRUCTURE_NAME_PREFIX = "LC_";
private MemberEnvAllowList() {
}
/**
* Derive the allowed NAME set from the given profiles plus {@link #INFRASTRUCTURE_PASSTHROUGH}.
* Deterministic (sorted) so generated scrub files are diffable run-to-run.
*/
public static Set<String> derive(Collection<FleetConfig.Profile> profiles) {
Set<String> derived = new TreeSet<>(INFRASTRUCTURE_PASSTHROUGH);
if (profiles != null) {
for (FleetConfig.Profile p : profiles) {
addIfPresent(derived, p.gitTokenEnv());
addIfPresent(derived, p.gitHostEnv());
addIfPresent(derived, p.tokenEnv());
if (p.env() != null) {
derived.addAll(p.env().keySet());
}
}
}
return Set.copyOf(derived);
}
/**
* Whether {@code name} survives the scrub when {@code allowedNames} is the derived set: an exact
* match, or an infrastructure-prefixed name ({@code LC_*}). Prefix rules live ONLY here and in
* the generated script's {@code case} pattern, which is written from this constant's value.
*/
public static boolean keeps(Set<String> allowedNames, String name) {
return allowedNames.contains(name) || name.startsWith(INFRASTRUCTURE_NAME_PREFIX);
}
/** The prefix rule as a zsh {@code case} pattern, so the script and Java cannot drift apart. */
public static String zshCasePattern() {
return INFRASTRUCTURE_NAME_PREFIX + "*";
}
private static void addIfPresent(Set<String> into, String name) {
if (name != null && !name.isBlank()) {
into.add(name);
}
}
}
@@ -1,15 +1,16 @@
package dev.ltms.bridged.member;
package dev.ltms.fleet.member;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.node.ObjectNode;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.WorkspaceControl;
import dev.ltms.bridged.peer.Capability;
import dev.ltms.bridged.peer.CharterReceipt;
import dev.ltms.bridged.peer.PeerHandle;
import dev.ltms.bridged.peer.SpawnRequest;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.WorkspaceControl;
import dev.ltms.fleet.peer.Capability;
import dev.ltms.fleet.peer.CharterReceipt;
import dev.ltms.fleet.peer.PeerHandle;
import dev.ltms.fleet.peer.PeerLauncher;
import dev.ltms.fleet.peer.SpawnRequest;
import java.io.IOException;
import java.io.UncheckedIOException;
@@ -23,6 +24,9 @@ import java.util.function.Function;
import java.util.function.LongSupplier;
import java.util.function.Supplier;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* The {@link HerdrPeerLauncher} adapter for <strong>opencode</strong> — an open-source,
* provider-agnostic terminal coding agent. Its whole reason for existing is to prove the
@@ -33,7 +37,7 @@ import java.util.function.Supplier;
* <p>The divergences from {@link ClaudeCodeLauncher}, all confined to {@link #buildLaunch}:
* <ul>
* <li><strong>No subscription boundary.</strong> opencode carries no {@code ANTHROPIC_BASE_URL}
* and there is no {@link dev.ltms.bridged.guard.SubscriptionGuard} — the guard is a
* and there is no {@link dev.ltms.fleet.guard.SubscriptionGuard} — the guard is a
* Claude-private concern, not part of the SPI. opencode reads the operator's own provider
* credentials from its global {@code auth.json}; the bridge injects none.</li>
* <li><strong>File-based MCP mount + instructions.</strong> opencode has no inline
@@ -49,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";
@@ -70,7 +76,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* matches the legacy non-blocking spawn semantics. Config dirs are created under the JVM temp dir.
*/
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env) {
this(agents, spaces, profiles, defaultProfile, env, 0,
System::currentTimeMillis, () -> sleepUninterruptibly(300),
@@ -82,7 +88,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* the pane reports an injectable state or {@code spawnReadyTimeoutMs} elapses.
*/
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs) {
this(agents, spaces, profiles, defaultProfile, env,
@@ -95,10 +101,10 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* <em>role</em>; a profile may still override it with its own {@code tabLabel}.
*/
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<BridgedConfig.Fleet> fleet) {
Supplier<FleetConfig.Fleet> fleet) {
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
defaultConfigRoot(), defaultDiscoveryRoot(), fleet);
@@ -108,11 +114,11 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* Production constructor, plus the CB-596 {@code memberCredentials} policy supplier.
*/
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs, long spawnReadyPollMs,
Supplier<BridgedConfig.Fleet> fleet,
Supplier<BridgedConfig.MemberCredentials> memberCredentials) {
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials);
@@ -138,7 +144,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* {@link OpenCodeSessionDiscovery})
*/
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
@@ -153,12 +159,12 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* @param fleet live fleet config, read once for each spawn
*/
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Path configRoot, Path discoveryRoot,
Supplier<BridgedConfig.Fleet> fleet) {
Supplier<FleetConfig.Fleet> fleet) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet);
this.configRoot = configRoot;
@@ -169,13 +175,13 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* Full testability constructor, plus the CB-596 {@code memberCredentials} policy supplier.
*/
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
Map<String, FleetConfig.Profile> profiles, String defaultProfile,
Function<String, String> env,
long spawnReadyTimeoutMs,
LongSupplier nowMillis, Runnable sleeper,
Path configRoot, Path discoveryRoot,
Supplier<BridgedConfig.Fleet> fleet,
Supplier<BridgedConfig.MemberCredentials> memberCredentials) {
Supplier<FleetConfig.Fleet> fleet,
Supplier<FleetConfig.MemberCredentials> memberCredentials) {
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
this.configRoot = configRoot;
@@ -201,11 +207,22 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* grant; and select the model with {@code -m}.
*/
@Override
protected Launch buildLaunch(BridgedConfig.Profile cfg, LaunchSpec spec) {
protected Launch buildLaunch(FleetConfig.Profile cfg, LaunchSpec spec) {
Map<String, String> workerEnv = baseEnv(cfg);
// A config file is needed for the bridge MCP mount, a member charter, or a pinned endpoint (CB-508).
if (cfg.hasMcp() || spec.charter() != null || hasCustomProvider(cfg)) {
workerEnv.put("OPENCODE_CONFIG", writeConfig(cfg, spec.charter()).toString());
// 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 \"<provider>/<model>\" "
+ "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<String> argv = argvWithResume(argvWithModel(argvWithAuto(cfg), cfg), spec.resumeSessionId());
@@ -238,7 +255,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* primary's Anthropic subscription, and an opencode process has no Anthropic credential path
* at all. Pointing it at a local vLLM cannot leak the subscription.
*/
private static boolean hasCustomProvider(BridgedConfig.Profile cfg) {
private static boolean hasCustomProvider(FleetConfig.Profile cfg) {
return cfg.baseUrl() != null && !cfg.baseUrl().isBlank();
}
@@ -252,7 +269,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* own git worktree on its own branch, is off-subscription, and cannot merge — the lead is the
* gate.
*/
private List<String> argvWithAuto(BridgedConfig.Profile cfg) {
private List<String> argvWithAuto(FleetConfig.Profile cfg) {
List<String> argv = mutableArgv(cfg.argv());
argv.add("--auto");
return argv;
@@ -275,7 +292,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
}
/** The launch argv plus, when a model is configured, the opencode {@code -m provider/model} flag. */
private List<String> argvWithModel(List<String> argv, BridgedConfig.Profile cfg) {
private List<String> argvWithModel(List<String> argv, FleetConfig.Profile cfg) {
if (cfg.model() != null && !cfg.model().isBlank()) {
argv.add("-m");
argv.add(cfg.model());
@@ -289,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(BridgedConfig.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();
@@ -320,15 +337,42 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
root.putArray("instructions").add(charter.toAbsolutePath().toString());
}
if (cfg.hasMcp()) {
ObjectNode bridge = root.putObject("mcp").putObject("bridge");
bridge.put("type", "remote");
bridge.put("url", cfg.mcpUrl());
bridge.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
@@ -349,7 +393,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* <p>The provider id comes from the {@code provider/model} selector in {@code model:}, so one
* field drives both the declaration and the {@code -m} flag and they cannot drift apart.
*/
private void addCustomProvider(ObjectNode root, BridgedConfig.Profile cfg) {
private void addCustomProvider(ObjectNode root, FleetConfig.Profile cfg) {
String[] parts = splitModelSelector(cfg);
String providerId = parts[0];
String modelId = parts[1];
@@ -362,7 +406,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
options.put("baseURL", openAiBaseUrl(cfg.baseUrl()));
// vLLM and friends usually ignore the key, but the AI SDK still requires a non-empty one.
String token = resolveEnv(cfg.tokenEnv());
options.put("apiKey", (token == null || token.isBlank()) ? "bridged-local-noauth" : token);
options.put("apiKey", (token == null || token.isBlank()) ? "fleetd-local-noauth" : token);
provider.putObject("models").putObject(modelId).put("name", modelId);
}
@@ -372,19 +416,69 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
* needs both, so a bare model name is rejected loudly rather than silently falling back to the
* default gateway — a worker quietly talking to the wrong endpoint is the failure this avoids.
*/
private static String[] splitModelSelector(BridgedConfig.Profile cfg) {
String model = cfg.model();
int slash = model == null ? -1 : model.indexOf('/');
if (model == null || model.isBlank() || slash <= 0 || slash == model.length() - 1) {
private static String[] splitModelSelector(FleetConfig.Profile cfg) {
String[] parts = splitProviderModel(cfg.model());
if (parts == null) {
throw new IllegalArgumentException(
"profile " + cfg.profile() + " sets baseUrl (a pinned opencode endpoint) so"
+ " model: must be \"<provider>/<model>\", 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.
*
* <p>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.<p>.models.<m>.limit.context}, which bounds the window opencode compacts
* <em>within</em> rather than compacting exactly AT it the way Claude Code's {@code --autocompact}
* does.
*
* <p>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.<p>.models.<m>.limit} override, which opencode merges over its own built-in
* provider definition.
*
* <p>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}).
*
* <p>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,
@@ -496,7 +590,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
/** Whether any configured profile opts into a git-forge token (required for {@link Capability#SELF_PR}). */
private boolean hasGitTokenProfile() {
return profileConfigs().stream().anyMatch(BridgedConfig.Profile::hasGitToken);
return profileConfigs().stream().anyMatch(FleetConfig.Profile::hasGitToken);
}
// --- CB-117 reap predicate (opencode prefix), kept for direct unit testing -----------------
@@ -1,4 +1,4 @@
package dev.ltms.bridged.member;
package dev.ltms.fleet.member;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
@@ -1,8 +1,8 @@
package dev.ltms.bridged.metrics;
package dev.ltms.fleet.metrics;
import dev.ltms.bridged.msg.ReplyInbox;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.fleet.msg.ReplyInbox;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.session.MemberSession;
import java.util.LinkedHashMap;
import java.util.Map;
@@ -13,33 +13,33 @@ import java.util.Map;
*
* <p>The set is deliberately small: each series maps to a failure mode this project has actually
* hit, not to whatever was easy to count. The two worth watching in practice are
* {@code bridged_sends_total{outcome="completion_fallback"}} — a rising share means turn detection
* {@code fleet_sends_total{outcome="completion_fallback"}} — a rising share means turn detection
* is degrading, the CB-115/116/118 failure family — and
* {@code bridged_push_nudges_total{outcome="exhausted"}}, which means the primary stopped draining
* {@code fleet_push_nudges_total{outcome="exhausted"}}, which means the primary stopped draining
* its inbox and CB-307's active push gave up.
*/
public final class BridgedMetrics {
public final class FleetMetrics {
/** Counter: delegated sends by terminal outcome. */
public static final String SENDS = "bridged_sends_total";
public static final String SENDS = "fleet_sends_total";
/** Counter: worker replies by the path that carried them (rendezvous vs stranded-to-inbox). */
public static final String REPLIES = "bridged_replies_total";
public static final String REPLIES = "fleet_replies_total";
/** Counter: push-loop nudges to the primary, by outcome. */
public static final String PUSH_NUDGES = "bridged_push_nudges_total";
public static final String PUSH_NUDGES = "fleet_push_nudges_total";
/** Counter: idle-lead heartbeat nudges to the lead, by outcome (CB-551). */
public static final String HEARTBEAT_NUDGES = "bridged_lead_heartbeat_nudges_total";
public static final String HEARTBEAT_NUDGES = "fleet_lead_heartbeat_nudges_total";
/** Counter: spawn attempts by peer kind and outcome. */
public static final String SPAWNS = "bridged_spawns_total";
public static final String SPAWNS = "fleet_spawns_total";
/** Counter: herdr socket calls by method and outcome. */
public static final String HERDR_CALLS = "bridged_herdr_calls_total";
public static final String HERDR_CALLS = "fleet_herdr_calls_total";
/** Counter: rejected requests by reason (CB-501). */
public static final String AUTH_FAILURES = "bridged_auth_failures_total";
public static final String AUTH_FAILURES = "fleet_auth_failures_total";
/** Gauge: session census by lifecycle state. */
public static final String SESSIONS = "bridged_sessions";
public static final String SESSIONS = "fleet_sessions";
/** Gauge: undrained replies held per target. */
public static final String INBOX_DEPTH = "bridged_inbox_depth";
public static final String INBOX_DEPTH = "fleet_inbox_depth";
private BridgedMetrics() {
private FleetMetrics() {
}
/**
@@ -1,4 +1,4 @@
package dev.ltms.bridged.metrics;
package dev.ltms.fleet.metrics;
import java.util.Map;
import java.util.NavigableMap;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.msg;
package dev.ltms.fleet.msg;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.msg;
package dev.ltms.fleet.msg;
import java.util.LinkedHashMap;
import java.util.List;
@@ -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.
*
* <p>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.
*
* <p>Note what is <em>not</em> 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<LeadMessage> 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();
}
@@ -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.
*
* <p>{@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.
*
* <p><strong>Status-gated, exactly like {@link ReplyPushLoop}.</strong> 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.
*
* <p><strong>Ack only after delivery.</strong> 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.
*
* <p><strong>One message per tick.</strong> 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 <em>before</em> 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<Map<String, String>> 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<Map<String, String>> 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<LeadMessage> 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:
*
* <ol>
* <li>a lead whose configured name equals {@code selfCoordId} — the explicit, unambiguous case;</li>
* <li>otherwise the sole lead, when this daemon recognises exactly one;</li>
* <li>otherwise nothing, and the message waits.</li>
* </ol>
*
* <p>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<String, String> 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;
}
}
@@ -1,11 +1,11 @@
package dev.ltms.bridged.msg;
package dev.ltms.fleet.msg;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.session.MemberSession;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -99,7 +99,7 @@ public final class LeadHeartbeatLoop {
/** Count one nudge outcome when a registry is wired; a no-op in unit tests. */
private void countNudge(String outcome) {
if (metrics != null) {
metrics.inc(BridgedMetrics.HEARTBEAT_NUDGES, "outcome", outcome);
metrics.inc(FleetMetrics.HEARTBEAT_NUDGES, "outcome", outcome);
}
}
@@ -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.
*
* <p><strong>Single-target, unlike {@link AmqpReplyInbox}.</strong> {@code AmqpReplyInbox}
* multiplexes many workers' reply queues under one gateway connection. A {@code LeadMailbox}
* instance is simpler: it owns exactly <em>one</em> queue — this daemon's own
* {@code lead.<selfCoordId>.inbox} — declared and consumed the moment it is constructed. There is
* no {@code own}/{@code release} pair to call separately; a daemon either runs a {@code LeadMailbox}
* for its own coord-id, or it does not run one at all.
*
* <p><strong>Consume-and-hold with deferred manual ack</strong> — same mapping as
* {@code AmqpReplyInbox}. The constructor declares the durable queue and starts a manual-ack
* consumer that pulls persistent messages into an in-memory {@code held} map (keyed by
* {@link LeadMessage#msgId()}) but does not ack them. {@link #peek} returns a non-destructive
* snapshot; {@link #ack} acks the broker delivery-tag and drops the entry. A message that is never
* acked (a crash, a bounce) survives — the broker redelivers it to the next connection that owns
* the queue.
*
* <p><strong>Publishing does not imply owning.</strong> {@link #publish} sends to
* {@code lead.<toCoordId>.inbox} over a dedicated confirm-mode channel; it never declares that
* queue as owned and never attaches a consumer to it. A sender that has never opened its own
* {@code LeadMailbox} for {@code toCoordId} can still publish to it, exactly as CB-308 federation
* requires. Publish blocks for the broker's publisher confirm (persistent delivery, {@code
* mandatory=true}) and throws {@link IllegalStateException} on an unroutable return, a nack, or a
* timeout — the caller must not report success for a black-holed message.
*
* <p><strong>Recovery.</strong> The connection is opened with automatic + topology recovery
* enabled, mirroring {@code AmqpReplyInbox}: on reconnect the broker hands out fresh delivery tags,
* so the held snapshot is cleared (dedup by {@code msgId} still prevents any double-queue on
* redelivery) and any publish still awaiting its confirm is failed rather than left to idle out
* the confirm timeout against a sequence number that means nothing on the new channel.
*/
public final class LeadMailbox implements 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<String, Held> held = new LinkedHashMap<>();
/**
* A dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume + ack)
* so a publish confirm round trip never blocks under {@link #channelLock} and stalls an ack.
*/
private final Channel publishChannel;
private final Object publishChannelLock = new Object();
/** In-flight publishes awaiting their confirm, keyed by the publish channel's sequence number. */
private final ConcurrentSkipListMap<Long, Pending> pendingBySeq = new ConcurrentSkipListMap<>();
/** The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery tag. */
private final ConcurrentHashMap<String, Pending> pendingByMsgId = new ConcurrentHashMap<>();
/** A message pulled off the broker but not yet acked: its delivery-tag plus the deserialized envelope. */
private record Held(long deliveryTag, LeadMessage message) {}
/** A publish awaiting its confirm; {@link #returned} records whether the broker already returned it. */
private static final class Pending {
final String msgId;
final CompletableFuture<Void> confirmed = new CompletableFuture<>();
volatile boolean returned;
Pending(String msgId) {
this.msgId = msgId;
}
}
/**
* Connect to {@code uri} (the shared cross-host coordination vhost, e.g.
* {@code amqp://guest:guest@127.0.0.1:5672/coord}) and own {@code selfCoordId}'s mailbox, with
* {@link #DEFAULT_PREFETCH}.
*/
public static LeadMailbox open(String uri, String selfCoordId) {
return open(uri, selfCoordId, DEFAULT_PREFETCH);
}
/** As {@link #open(String, String)}, with an explicit consumer prefetch. */
public static LeadMailbox open(String uri, String selfCoordId, int prefetch) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares the queue and re-attaches the consumer.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
return new LeadMailbox(factory.newConnection("bridged-lead-mailbox"), selfCoordId, prefetch);
} catch (Exception e) {
throw new IllegalStateException("cannot connect to AMQP coordination broker at " + uri, e);
}
}
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for tests). */
LeadMailbox(Connection connection, String selfCoordId) {
this(connection, selfCoordId, DEFAULT_PREFETCH);
}
/** As above, with an explicit prefetch (injection seam for tests). */
LeadMailbox(Connection connection, String selfCoordId, int prefetch) {
this.connection = connection;
this.selfCoordId = selfCoordId;
try {
this.channel = connection.createChannel();
// Bound the held backlog — must be set before basicConsume.
this.channel.basicQos(prefetch);
this.publishChannel = connection.createChannel();
this.publishChannel.confirmSelect();
this.publishChannel.addReturnListener(this::onReturn);
this.publishChannel.addConfirmListener(this::onAck, this::onNack);
own();
} catch (IOException e) {
throw new IllegalStateException("cannot open AMQP channel", e);
}
// On automatic recovery the broker redelivers unacked messages with FRESH delivery-tags; the
// tags we were holding are now stale. Drop the held snapshot so the re-attached consumer
// repopulates it with valid tags (dedup by msgId still prevents any double-queue). Any publish
// confirm still in flight when the connection dropped is equally stale — fail it now rather
// than let it silently ride out CONFIRM_TIMEOUT_MS.
if (connection instanceof Recoverable recoverable) {
recoverable.addRecoveryListener(new RecoveryListener() {
@Override
public void handleRecovery(Recoverable recoverable) {
synchronized (held) {
held.clear();
}
failPendingPublishesOnRecovery();
log.info("AMQP lead mailbox connection recovered; cleared held messages for fresh redelivery");
}
@Override
public void handleRecoveryStarted(Recoverable recoverable) {
// no-op: we act once recovery completes
}
});
}
}
/** Declare + consume this daemon's own {@code lead.<selfCoordId>.inbox}. Called once, at construction. */
private void own() throws IOException {
String queue = queueName(selfCoordId);
synchronized (channelLock) {
channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
channel.basicConsume(queue, false, deliverCallback(), _ -> { }); // autoAck=false: manual ack
}
log.debug("lead mailbox owns queue {} for coord-id {}", queue, selfCoordId);
}
/**
* Publish {@code msg} to {@code toCoordId}'s mailbox and block until the broker's publisher
* confirm for it lands. Does <em>not</em> imply owning or consuming {@code toCoordId}'s queue.
* Throws {@link IllegalStateException} if the message is returned as unroutable, nacked, or not
* confirmed within {@link #CONFIRM_TIMEOUT_MS} — the caller must treat that as a failed publish,
* not a lost-and-forgotten one.
*/
@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<LeadMessage> peek() {
synchronized (held) {
return held.values().stream().map(Held::message).toList();
}
}
/** Convenience: {@link #peek} the current snapshot, then {@link #ack} every message in it. */
public List<LeadMessage> drain() {
List<LeadMessage> snapshot = peek();
snapshot.forEach(m -> ack(m.msgId()));
return snapshot;
}
/** Remove the held message {@code msgId} and ack it on the broker. No-op if not held. */
@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<Long, Pending> covered = multiple
? pendingBySeq.headMap(seq, true)
: pendingBySeq.subMap(seq, true, seq, true);
for (var it = covered.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
if (ack && !pending.returned) {
pending.confirmed.complete(null);
} else if (ack) {
pending.confirmed.completeExceptionally(new IllegalStateException(
"lead message " + pending.msgId + " was returned as unroutable (mailbox not owned)"));
} else {
pending.confirmed.completeExceptionally(new IllegalStateException(
"broker nacked publish of lead message " + pending.msgId));
}
}
}
/**
* Fail every publish still awaiting its confirm — their sequence numbers are stale after
* recovery. Guarded by {@link #publishChannelLock}, the same lock {@link #publish} holds while
* it takes its sequence number and registers its {@link Pending} — see
* {@code AmqpReplyInbox.failPendingPublishesOnRecovery}'s javadoc for the full race analysis this
* mirrors. Package-private only so a unit test can drive it directly without a live broker
* reconnect.
*/
void failPendingPublishesOnRecovery() {
synchronized (publishChannelLock) {
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
pending.confirmed.completeExceptionally(new IllegalStateException(
"AMQP connection recovered mid-publish; confirm status of lead message "
+ pending.msgId + " is unknown"));
}
}
}
/**
* Fail every publish still awaiting its confirm with a clear, immediate error instead of leaving
* it to time out after {@link #CONFIRM_TIMEOUT_MS} once the channels are closed underneath it.
*/
private void failPendingPublishesOnClose() {
synchronized (publishChannelLock) {
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
pending.confirmed.completeExceptionally(new IllegalStateException(
"lead mailbox closed while publish of lead message " + pending.msgId
+ " was still awaiting its confirm"));
}
}
}
/** The durable queue name a coord-id's mailbox lives on: {@code lead.<coordId>.inbox}. */
public static String queueName(String coordId) {
return QUEUE_PREFIX + coordId + QUEUE_SUFFIX;
}
@Override
public void close() {
failPendingPublishesOnClose();
try {
channel.close();
} catch (Exception e) {
log.debug("AMQP lead mailbox channel close: {}", e.toString());
}
try {
publishChannel.close();
} catch (Exception e) {
log.debug("AMQP lead mailbox publish channel close: {}", e.toString());
}
try {
connection.close();
} catch (Exception e) {
log.debug("AMQP lead mailbox connection close: {}", e.toString());
}
}
}
@@ -0,0 +1,22 @@
package dev.ltms.fleet.msg;
/**
* Wire envelope for a lead-to-lead message carried over {@link LeadMailbox}.
*
* <p>Unlike {@link ReplyInbox.InboxMessage} (a worker→primary reply, addressed only by the single
* gateway that owns the worker), a lead message crosses independently-owned daemons — possibly on
* different hosts — so it carries an explicit sender ({@code from}) as well as the recipient
* ({@code to}): the recipient needs the sender's coord-id to reply back.
*
* <p>{@code from} and {@code to} are globally-unique lead coordination ids (e.g. {@code "mac-opus"},
* {@code "fleet01-lead"}) — NOT herdr terminal ids. A herdr terminal id is meaningful only on the
* host that owns it, so it cannot address a lead running on another daemon; a coord-id is chosen
* by configuration ({@code coordinator.selfId}) precisely so it means the same thing everywhere.
*
* @param msgId idempotency id; a redelivered duplicate (at-least-once delivery) is deduped on this
* @param from the sending lead's coord-id
* @param to the receiving lead's coord-id — identifies the mailbox this message is held on
* @param content the message text
*/
public record LeadMessage(String msgId, String from, String to, String content) {
}
@@ -1,10 +1,10 @@
package dev.ltms.bridged.msg;
package dev.ltms.fleet.msg;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.inject.Injector;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -284,13 +284,13 @@ public final class MessageService {
*/
public boolean reply(String session, String content) {
if (rendezvous.resolve(session, content)) {
count(BridgedMetrics.REPLIES, "path", "rendezvous");
count(FleetMetrics.REPLIES, "path", "rendezvous");
return true; // a live send took it — unchanged fast path
}
inbox.publish(session, UUID.randomUUID().toString(), content);
// A rising inbox share is the signal CB-307 exists to make visible: the worker finished but
// nobody was waiting, so delivery now depends on the push loop and a drain.
count(BridgedMetrics.REPLIES, "path", "inbox");
count(FleetMetrics.REPLIES, "path", "inbox");
if (pushLoop != null) {
pushLoop.onReplyQueued(session);
}
@@ -308,7 +308,7 @@ public final class MessageService {
private Reply recorded(Reply r) {
String label = sendOutcomeLabel(r.outcome());
if (label != null) {
count(BridgedMetrics.SENDS, "outcome", label);
count(FleetMetrics.SENDS, "outcome", label);
}
return r;
}
@@ -512,7 +512,7 @@ public final class MessageService {
// CB-582: the question just became visible via fleet_poll (Phase.ASKING) for an async
// (wait:false) delegation — nudge the lead's own pane the same way a terminal ticket does
// (CB-588), since the lead's normal poll cadence is minutes away and the reverse-rendezvous
// window (~55s, see BridgeMcp/BridgedApp) is far shorter. A blocking (wait:true) send has
// window (~55s, see FleetMcp/FleetApp) is far shorter. A blocking (wait:true) send has
// no Task and gets the question directly in its own reply, so task == null there — nothing
// to nudge.
if (task != null && pushLoop != null) {
@@ -1,4 +1,4 @@
package dev.ltms.bridged.msg;
package dev.ltms.fleet.msg;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.msg;
package dev.ltms.fleet.msg;
import java.util.List;
@@ -1,10 +1,10 @@
package dev.ltms.bridged.msg;
package dev.ltms.fleet.msg;
import dev.ltms.bridged.herdr.AgentControl;
import dev.ltms.bridged.herdr.AgentStatus;
import dev.ltms.bridged.mcp.PrimaryRegistry;
import dev.ltms.bridged.metrics.BridgedMetrics;
import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.fleet.herdr.AgentControl;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.mcp.PrimaryRegistry;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -118,7 +118,7 @@ public final class ReplyPushLoop {
/** Count one nudge outcome when a registry is wired; a no-op in unit tests. */
private void countNudge(String outcome) {
if (metrics != null) {
metrics.inc(BridgedMetrics.PUSH_NUDGES, "outcome", outcome);
metrics.inc(FleetMetrics.PUSH_NUDGES, "outcome", outcome);
}
}
@@ -364,7 +364,7 @@ public final class ReplyPushLoop {
/**
* Called when an async ticket's worker pauses mid-turn in {@code fleet_ask} (CB-582): the
* question is now visible via {@code fleet_poll} (Phase.ASKING), but the reverse-rendezvous
* window it opened with (~55s default, see {@code BridgeMcp}/{@code BridgedApp}) is far shorter
* window it opened with (~55s default, see {@code FleetMcp}/{@code FleetApp}) is far shorter
* than a lead's normal minutes-long poll cadence — exactly the gap this closes. Resolves the
* delegating lead the same way {@link #onTicketTerminal} does and coalesces onto the same
* per-lead schedule (CB-590).
@@ -1,4 +1,4 @@
package dev.ltms.bridged.msg;
package dev.ltms.fleet.msg;
import java.util.concurrent.CompletableFuture;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.peer;
package dev.ltms.fleet.peer;
/**
* Declared capabilities of a {@link PeerLauncher}. The protocol is the union across all
@@ -1,4 +1,4 @@
package dev.ltms.bridged.peer;
package dev.ltms.fleet.peer;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.peer;
package dev.ltms.fleet.peer;
import java.util.Locale;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.peer;
package dev.ltms.fleet.peer;
/**
* An opaque handle returned by {@link PeerLauncher#spawn(SpawnRequest)}. The core routes on
@@ -0,0 +1,195 @@
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
* config; today's single adapter is the {@code ClaudeCodeLauncher} / Claude Code over herdr.
*
* <p>The core delegates spawn and teardown to this interface without knowing how the peer is set
* up. Environment variables, CLI flags, subscription guards, transport (herdr tab/pane) layout,
* and naming conventions are all adapter-private — the core sees only the returned
* {@link PeerHandle} whose {@code id()} is the registry/routing key.
*
* <p>The interface is a superset of what {@code SessionManager} and {@code Fleetd.main} call
* on the concrete launcher today.
*/
public interface PeerLauncher {
/**
* The name every launcher gives the bridge's MCP server in the config it writes for its peer.
* The peer's tools are addressed as {@code mcp__<this>__fleet_*}, and {@code CLAUDE.md}'s
* role-detection ladder names that prefix, so the two must agree.
*
* <p>It is a constant because three launchers write it — {@code ClaudeCodeLauncher} and
* {@code LeadLauncher} into a {@code --mcp-config} literal, {@code OpenCodeLauncher} into an
* {@code opencode.json} node. Three hand-written copies of one name is how a rename lands in
* two of them (CB-632).
*/
String MCP_MOUNT_NAME = "fleet";
/**
* Shared IDE-guidance text (CB-634), delivered per-backend as an on-disk overlay rather than
* any one adapter's system-prompt charter, so a project's own {@code CLAUDE.md} is never
* clobbered. It pins every {@code ide_*} call to the member's own worktree, which is the whole
* point of the mechanism. Both launchers render their own overlay from this single source.
*
* @param projectPath the path the member must pin every {@code ide_*} call to — the module dir
* IntelliJ opened as the project, which is {@link #ideProjectPath} of the
* member's own worktree (the worktree root when no module subdir is set)
*/
static String ideOverlayText(String projectPath) {
return "## IDE code intelligence — your worktree only\n"
+ "An IntelliJ IDE Index MCP server is mounted as `mcp__intellij__ide_*`. Prefer it "
+ "over `grep`/`find` for symbol lookups, references, call and type hierarchy, and "
+ "diagnostics — it resolves the real AST, text search does not.\n\n"
+ "Every `ide_*` call MUST pass `project_path: \"" + projectPath + "\"` — your own "
+ "worktree — and never any other path. A call without it errors "
+ "`multiple_projects_open`; a call with a different path reads another checkout, "
+ "not your changes. This is not the primary's IDE: it is your worktree, pinned to "
+ "you.";
}
/**
* The absolute path IntelliJ must open as the project, and the {@code project_path} the overlay
* pins (CB-634). It is {@code cwd} resolved against {@code ideProjectDir}. The distinction
* matters because this repo (like {@code fleet/fleetd}) keeps its Maven module in a subdir
* ({@code bridged/}), not at the worktree root: opening the root imports no module and
* {@code ide_*} resolves nothing, so the module dir is the correct pin and open target.
*
* @param cwd the member's worktree root
* @param ideProjectDir repo-relative module subdir, or {@code null}/blank for the worktree root
* @return the absolute, normalized module dir as a string
*/
static String ideProjectPath(String cwd, String ideProjectDir) {
Path base = Path.of(cwd);
if (ideProjectDir == null || ideProjectDir.isBlank()) {
return base.toString();
}
return base.resolve(ideProjectDir).normalize().toString();
}
/**
* Best-effort: open {@code projectPath} in the host IDE by running {@code openCommand} with
* every {@code {dir}} replaced by {@code projectPath} (CB-634 auto-open). The command runs
* through {@code /bin/sh -c} so an operator can set env inline — e.g.
* {@code "env DISPLAY=:10.0 idea {dir}"} — because the daemon's own env may lack {@code DISPLAY}.
*
* <p>A blank command is a no-op: the profile opted into IDE MCP but not auto-open, so the
* operator opens the module by hand. The child process is detached and its exit is not awaited;
* any failure is logged and swallowed, because a member must spawn whether or not an IDE is
* running. There is no close half yet (CB-634 defers it): an opened module stays open until the
* operator closes it, and opening the same module again just refocuses it.
*
* @param projectPath the module dir to open (typically {@link #ideProjectPath})
* @param openCommand the host command template, with {@code {dir}} substituted; null/blank ⇒ no-op
* @param log the calling launcher's logger, for the best-effort WARN
*/
static void openInIde(String projectPath, String openCommand, Logger log) {
if (openCommand == null || openCommand.isBlank()) {
return;
}
String cmd = openCommand.replace("{dir}", projectPath);
try {
new ProcessBuilder("/bin/sh", "-c", cmd)
.redirectOutput(ProcessBuilder.Redirect.DISCARD)
.redirectError(ProcessBuilder.Redirect.DISCARD)
.start();
log.info("CB-634 auto-open: launched IDE open for {}", projectPath);
} catch (Exception e) {
log.warn("CB-634 auto-open of '{}' failed (member still spawns): {}", projectPath, e.getMessage());
}
}
/**
* The set of {@link Capability capabilities} this launcher declares. A peer whose profile
* opts into a git-forge token should include {@link Capability#SELF_PR}; the base set for
* the Claude Code herdr adapter is always {@code MID_TURN_ASK, WORKTREE, ORPHAN_REAP}.
*/
Set<Capability> capabilities();
/**
* The capabilities of the adapter that {@code profileName} resolves to (null/blank → the
* default profile, the same resolution {@link #spawn} uses). Distinct from {@link
* #capabilities()}, which unions every configured adapter: a caller that must know whether
* <em>this</em> profile's backend supports a capability — e.g. {@link Capability#SESSION_RESUME}
* before honoring {@link SpawnRequest#resumeSessionId()} — needs the per-profile answer, not
* the fleet-wide union, or a mixed fleet could OK a resume that lands on a non-supporting
* adapter (CB-584).
*
* @throws IllegalArgumentException if the profile is unknown and no default is configured
*/
Set<Capability> capabilitiesFor(String profileName);
/**
* {@code profileName}/requestedCwd null/blank → default resolution. Returns after the peer
* process is live (env + argv + placement complete). Never returns {@code null}.
*
* @param req the spawn parameters (profile, requested cwd, caller cwd)
* @return a handle whose {@link PeerHandle#id()} is the registry/routing key
* @throws IllegalArgumentException if the profile is unknown and no default is configured
*/
PeerHandle spawn(SpawnRequest req);
/**
* The configured worker profile names — the set of names {@code spawn(profileName)} accepts.
*/
Set<String> profiles();
/**
* The profile a no-argument {@link #spawn(SpawnRequest)} uses, or {@code null} if none is configured.
*/
String defaultProfile();
/**
* Resolve the effective working directory for a spawn {@code req} without actually spawning.
* Resolution order: requestedCwd → profile cwd → callerCwd → daemon cwd.
*
* @return the resolved absolute path, never null/blank
*/
String effectiveCwd(SpawnRequest req);
/**
* The parity-overlay file list for {@code profileName} (default list when unset). Used by
* worktree provisioning to copy config files into the isolated checkout before spawning.
*/
List<String> parityOverlay(String profileName);
/**
* The set of all agents this launcher currently tracks, transport-specific. Each element
* exposes at minimum a pane-like {@code id()} matching this launcher's {@link PeerHandle}
* scheme, plus transport-level status. Callers merge this set with the session registry to
* build a live roster view.
*/
List<?> list();
/**
* Reap orphaned peers left behind by a prior daemon process. Only peers whose naming scheme
* matches this launcher's and whose nonce differs from the current process are eligible.
* Best-effort: a failure to list or to stop any one peer is logged and never aborts startup.
*
* @return the number of orphaned peers reaped
*/
int reapOrphanWorkers();
/**
* Tear a peer down by its registry/routing key ({@link PeerHandle#id()}). Tolerates an
* already-gone peer. Also cleans up launcher-private resources (e.g. empty dedicated tabs)
* when safe to do so.
*/
void stop(String id);
/**
* Discard the context of the peer identified by {@code id}. Implementations must bypass normal
* bridge delivery/turn accounting. Unsupported peer kinds return {@code false} without sending
* a guessed command.
*
* @return {@code true} when a reset was sent and its status transition must settle before reuse
*/
boolean clearContext(String id);
}
@@ -1,4 +1,4 @@
package dev.ltms.bridged.peer;
package dev.ltms.fleet.peer;
/**
* Thrown when a {@link PeerLauncher} starts a peer process but the peer
@@ -1,4 +1,4 @@
package dev.ltms.bridged.peer;
package dev.ltms.fleet.peer;
/**
* What a caller asks for when spawning a peer.
@@ -1,4 +1,4 @@
package dev.ltms.bridged.placement;
package dev.ltms.fleet.placement;
import java.util.LinkedHashMap;
import java.util.Map;
@@ -8,7 +8,7 @@ import java.util.concurrent.ConcurrentHashMap;
import java.util.function.LongSupplier;
/**
* Where a credential (not a profile — see {@code BridgedConfig.Profile#effectiveCredentialId()})
* Where a credential (not a profile — see {@code FleetConfig.Profile#effectiveCredentialId()})
* sits out a cooldown after a {@code BACKEND_EXHAUSTED} classification (CB-578 stage B), so a fresh
* spawn does not walk straight back onto the account that just refused on a usage limit.
*
@@ -16,7 +16,7 @@ import java.util.function.LongSupplier;
* models on the same OpenAI account) share one quarantine — {@link #quarantine} one credential id
* and every profile whose {@code effectiveCredentialId()} equals it is quarantined too, without this
* class knowing anything about profiles at all. That mapping is the caller's job (see
* {@code CompositePeerLauncher} and {@code dev.ltms.bridged.inject.ExhaustionSink}).
* {@code CompositePeerLauncher} and {@code dev.ltms.fleet.inject.ExhaustionSink}).
*
* <p>The clock is injected ({@link LongSupplier}, conventionally {@code System::nanoTime} like
* {@code FleetHealthMonitor}), never read inline, so a quarantine's expiry is testable without a
@@ -1,4 +1,4 @@
package dev.ltms.bridged.placement;
package dev.ltms.fleet.placement;
/**
* Backward-compatible placement: an unqualified spawn always resolves to the configured default
@@ -1,4 +1,4 @@
package dev.ltms.bridged.placement;
package dev.ltms.fleet.placement;
/**
* A profile (and, in CB-308, a host) that can be chosen by a {@link PlacementPolicy}.
@@ -21,7 +21,7 @@ public record PlacementCandidate(String profile, String host, float weight, Inte
/**
* True when this candidate carries an explicit {@code weight <= 0} (CB-554) and must be
* skipped by every automatic policy — the same way a quarantined or unreachable candidate is
* skipped. {@code BridgedConfig.Profile}'s compact constructor already normalises "absent" to
* skipped. {@code FleetConfig.Profile}'s compact constructor already normalises "absent" to
* {@code 1.0} and "negative" to {@code 0.0}, so this is a plain threshold check here; it does
* not need to distinguish "explicit 0" from "absent" itself.
*
@@ -1,4 +1,4 @@
package dev.ltms.bridged.placement;
package dev.ltms.fleet.placement;
import java.util.List;
import java.util.Set;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.placement;
package dev.ltms.fleet.placement;
/**
* Thrown when a {@link PlacementPolicy} has no candidate available. Kept as a distinct type so
@@ -1,4 +1,4 @@
package dev.ltms.bridged.placement;
package dev.ltms.fleet.placement;
/**
* Factory for the built-in placement policies.
@@ -1,4 +1,4 @@
package dev.ltms.bridged.placement;
package dev.ltms.fleet.placement;
/**
* How {@code bridged} chooses a worker profile when a spawn names none. Implementations are
@@ -1,4 +1,4 @@
package dev.ltms.bridged.placement;
package dev.ltms.fleet.placement;
import java.util.ArrayList;
import java.util.List;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.placement;
package dev.ltms.fleet.placement;
import java.util.List;
import java.util.concurrent.atomic.AtomicInteger;
@@ -1,4 +1,4 @@
package dev.ltms.bridged.placement;
package dev.ltms.fleet.placement;
import java.util.List;
import java.util.Map;
@@ -1,25 +1,26 @@
package dev.ltms.bridged.rest;
package dev.ltms.fleet.rest;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import dev.ltms.bridged.auth.AuditLog;
import dev.ltms.bridged.auth.Authz;
import dev.ltms.bridged.auth.CallerResolver;
import dev.ltms.bridged.auth.Principal;
import dev.ltms.bridged.guard.GuardException;
import dev.ltms.bridged.metrics.Metrics;
import dev.ltms.bridged.herdr.Agent;
import dev.ltms.bridged.herdr.HerdrClient;
import dev.ltms.bridged.herdr.HerdrException;
import dev.ltms.bridged.inject.MemberPresence;
import dev.ltms.bridged.peer.PeerUnreachableException;
import dev.ltms.bridged.placement.PlacementException;
import dev.ltms.bridged.msg.MessageService;
import dev.ltms.bridged.session.SessionManager;
import dev.ltms.bridged.peer.MemberRole;
import dev.ltms.bridged.session.MemberSession;
import dev.ltms.bridged.session.WorktreeRequest;
import dev.ltms.bridged.peer.PeerLauncher;
import dev.ltms.fleet.auth.AuditLog;
import dev.ltms.fleet.auth.Authz;
import dev.ltms.fleet.auth.CallerResolver;
import dev.ltms.fleet.auth.Principal;
import dev.ltms.fleet.guard.GuardException;
import dev.ltms.fleet.metrics.FleetMetrics;
import dev.ltms.fleet.metrics.Metrics;
import dev.ltms.fleet.herdr.Agent;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.herdr.HerdrException;
import dev.ltms.fleet.inject.MemberPresence;
import dev.ltms.fleet.peer.PeerUnreachableException;
import dev.ltms.fleet.placement.PlacementException;
import dev.ltms.fleet.msg.MessageService;
import dev.ltms.fleet.session.SessionManager;
import dev.ltms.fleet.peer.MemberRole;
import dev.ltms.fleet.session.MemberSession;
import dev.ltms.fleet.session.WorktreeRequest;
import dev.ltms.fleet.peer.PeerLauncher;
import io.javalin.Javalin;
import io.javalin.http.Context;
import jakarta.servlet.http.HttpServlet;
@@ -41,7 +42,7 @@ import java.util.stream.Collectors;
* <p>Built from injected collaborators so tests supply fakes and run on an ephemeral
* port; {@code main} supplies the real Unix-socket client and worker service.
*/
public final class BridgedApp {
public final class FleetApp {
/** Default blocking window for a message; kept under typical HTTP idle timeouts. */
private static final long DEFAULT_MESSAGE_TIMEOUT_MS = 25_000;
@@ -68,7 +69,7 @@ public final class BridgedApp {
* behaved before CB-501. Retained so existing acceptance tests keep exercising handler
* behaviour without each needing an auth fixture.
*/
public BridgedApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet) {
this(herdr, workers, sessions, messages, presence, mcpServlet, null, null);
@@ -80,7 +81,7 @@ public final class BridgedApp {
* @param metrics registry to instrument and expose at {@code GET /metrics}; {@code null} omits
* the endpoint
*/
public BridgedApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
public FleetApp(HerdrClient herdr, PeerLauncher workers, SessionManager sessions,
MessageService messages, MemberPresence presence,
HttpServlet mcpServlet, CallerResolver auth, Metrics metrics) {
this.herdr = herdr;
@@ -104,7 +105,7 @@ public final class BridgedApp {
}
});
// CB-501: resolve identity once per request, before any handler. /mcp does NOT pass through
// here — it is a raw servlet on Jetty's context handler — so BridgeMcp enforces separately
// here — it is a raw servlet on Jetty's context handler — so FleetMcp enforces separately
// against the same CallerResolver. Any check that lives in only one place is not a control.
if (auth != null) {
app.before(ctx -> ctx.attribute(CALLER,
@@ -166,7 +167,7 @@ public final class BridgedApp {
private void countAuthFailure(String reason) {
if (metrics != null) {
metrics.inc("bridged_auth_failures_total", "reason", reason);
metrics.inc(FleetMetrics.AUTH_FAILURES, "reason", reason);
}
}
@@ -219,7 +220,7 @@ public final class BridgedApp {
return;
}
ctx.status(200).json(Map.of("agents",
workers.list().stream().map(Agent.class::cast).map(BridgedApp::view).toList()));
workers.list().stream().map(Agent.class::cast).map(FleetApp::view).toList()));
}
/** CB-304: bridge-owned roster merged with live herdr status by paneId. */
@@ -1,4 +1,4 @@
package dev.ltms.bridged.session;
package dev.ltms.fleet.session;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Some files were not shown because too many files have changed in this diff Show More