95a8dbcea9
#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.
789 lines
44 KiB
Java
789 lines
44 KiB
Java
package dev.ltms.fleet;
|
|
|
|
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.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;
|
|
import java.util.List;
|
|
import java.util.Map;
|
|
import java.util.Objects;
|
|
import java.util.Set;
|
|
import java.util.concurrent.Executors;
|
|
import java.util.concurrent.TimeUnit;
|
|
import java.util.concurrent.atomic.AtomicReference;
|
|
import java.util.function.Function;
|
|
import java.util.function.Predicate;
|
|
import java.util.function.Supplier;
|
|
import java.util.regex.Pattern;
|
|
import java.util.stream.Collectors;
|
|
|
|
/**
|
|
* {@code bridged} entry point. Wires the real herdr socket client to the REST app and
|
|
* 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 Fleetd {
|
|
|
|
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;
|
|
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 = 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.
|
|
reportRequiredSecrets(cfg);
|
|
// CB-596: an absent (or empty) memberCredentials: block blocks NOTHING — no credential
|
|
// name is hardcoded any more to fall back on. Say so loudly, the same way a missing
|
|
// secret is reported above, so upgrading past this commit never silently drops CB-592's
|
|
// protection.
|
|
reportMemberCredentialsGap(cfg);
|
|
// CB-559: `cfg` stays the startup snapshot — every validation and every piece of one-time
|
|
// wiring below reads it, and must, because those decisions cannot be unmade. `config` is the
|
|
// live reference the hot paths read per use. Which keys can actually move is ConfigRef's
|
|
// contract; adding a reader here does not make a key reloadable by itself.
|
|
ConfigRef config = new ConfigRef(configPath, cfg);
|
|
|
|
// The primary/host env that launched bridged must not be tainted.
|
|
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
|
guard.assertPrimaryClean(System.getenv());
|
|
|
|
// CB-501: refuse to start if the bind is wider than the auth mode can defend. Under
|
|
// loopback-trust, "not a known worker" means "the primary" — sound only because the OS
|
|
// refuses remote connections to a loopback socket. This throws rather than warns so the
|
|
// dangerous configuration cannot be reached by ignoring a log line.
|
|
cfg.validateAuthExposure();
|
|
cfg.validateLeadTabPrefixes();
|
|
// CB-542: a subscription:true profile whose env: reseats ANTHROPIC_BASE_URL/AUTH_TOKEN would
|
|
// reach an unguarded endpoint (the launcher skips SubscriptionGuard for it). Refuse at load.
|
|
cfg.validateSubscriptionProfiles();
|
|
cfg.validateCharters();
|
|
// CB-548: every architect slot must name a configured workers: profile — the strong-model
|
|
// backend the future spawn lifecycle would read. A stale reference dies here, not later.
|
|
cfg.validateMembers();
|
|
|
|
Path socket = cfg.herdrSocket() != null && !cfg.herdrSocket().isBlank()
|
|
? Path.of(cfg.herdrSocket())
|
|
: UnixSocketHerdrClient.defaultSocketPath();
|
|
|
|
UnixSocketHerdrClient herdr = UnixSocketHerdrClient.connect(socket, new com.fasterxml.jackson.databind.ObjectMapper());
|
|
|
|
AgentControl agents = new AgentControl(herdr);
|
|
WorkspaceControl spaces = new WorkspaceControl(herdr);
|
|
// 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, FleetConfig.Profile> claudeProfiles = new LinkedHashMap<>();
|
|
Map<String, FleetConfig.Profile> opencodeProfiles = new LinkedHashMap<>();
|
|
cfg.profiles().forEach((name, w) -> {
|
|
if (w.isOpenCode()) {
|
|
opencodeProfiles.put(name, w);
|
|
} else {
|
|
claudeProfiles.put(name, w);
|
|
}
|
|
});
|
|
List<HerdrPeerLauncher> adapters = new ArrayList<>();
|
|
// The claude-code adapter is the always-present default; keep it even with no profiles (so a
|
|
// bridge configured with no workers, or opencode-only, still has a well-defined base adapter)
|
|
// unless opencode is the only kind configured.
|
|
if (!claudeProfiles.isEmpty() || opencodeProfiles.isEmpty()) {
|
|
adapters.add(new ClaudeCodeLauncher(agents, spaces, guard,
|
|
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
|
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
|
() -> config.get().fleet(),
|
|
() -> config.get().memberCredentials()));
|
|
}
|
|
if (!opencodeProfiles.isEmpty()) {
|
|
adapters.add(new OpenCodeLauncher(agents, spaces,
|
|
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
|
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
|
() -> config.get().fleet(),
|
|
() -> config.get().memberCredentials()));
|
|
}
|
|
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 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()));
|
|
PeerLauncher workers = new CompositePeerLauncher(
|
|
adapters,
|
|
cfg.effectiveDefaultProfile(),
|
|
config,
|
|
profileName -> liveCountRef.get().apply(profileName),
|
|
quarantine);
|
|
// CB-504: under supervision (launchd/systemd) bridged can start before herdr's socket
|
|
// exists. The client itself is lazy — it connects per call — but the orphan reap below is
|
|
// the first thing that actually talks to herdr, so without this wait a boot-order race
|
|
// would crash the daemon into a restart loop. Wait, then degrade rather than die: serving
|
|
// with /healthz reporting "degraded" is strictly more useful than exiting.
|
|
boolean herdrUp = awaitHerdr(herdr);
|
|
if (herdrUp) {
|
|
// CB-117: herdr keeps worker panes alive across a daemon restart, and their ids died
|
|
// with the previous process — reap those leaked orphans now, before we start serving.
|
|
workers.reapOrphanWorkers();
|
|
} else {
|
|
log.warn("herdr did not answer within {}s — starting anyway; /healthz will report "
|
|
+ "degraded until it comes up. Orphaned worker panes (if any) were NOT reaped.",
|
|
HERDR_WAIT_SECONDS);
|
|
}
|
|
|
|
// CB-301: authoritative session registry + lifecycle FSM on top of ClaudeCodeLauncher.
|
|
// CB-301-ext: worktree provisioning seam, optionally rooted at a configured directory.
|
|
// CB-303 part 2: context cap is opt-in and disabled (0) when absent/null.
|
|
int contextCap = 0;
|
|
if (cfg.lifecycle() != null && cfg.lifecycle().contextCap() != null
|
|
&& cfg.lifecycle().contextCap() > 0) {
|
|
contextCap = cfg.lifecycle().contextCap();
|
|
}
|
|
boolean clearAfterTurn = cfg.lifecycle() != null && cfg.lifecycle().clearAfterTurn();
|
|
SessionManager sessions = new SessionManager(workers, new GitWorktrees(cfg.worktreeRoot()),
|
|
System::nanoTime, contextCap, clearAfterTurn);
|
|
liveCountRef.set(profileName -> (int) sessions.roster().stream()
|
|
.filter(s -> profileName.equals(s.profile()))
|
|
.count());
|
|
|
|
// CB-303 part 1: idle-ttl reaper — only when configured, defaults to disabled.
|
|
final SessionReaper reaper;
|
|
if (cfg.lifecycle() != null
|
|
&& cfg.lifecycle().idleTtlSeconds() != null
|
|
&& cfg.lifecycle().idleTtlSeconds() > 0) {
|
|
reaper = new SessionReaper(sessions, cfg.lifecycle().idleTtlSeconds());
|
|
reaper.start();
|
|
} else {
|
|
reaper = null;
|
|
}
|
|
|
|
// CB-530: every pane the config names as a lead, merged from `leaders:` and the legacy
|
|
// singular pin. PrimaryRegistry below still tracks ONE terminal — it addresses the push
|
|
// loop's nudges, which need a single destination — so it keeps the legacy pin.
|
|
Map<String, String> leadTerminals = cfg.leaderTerminals();
|
|
if (leadTerminals.size() > 1) {
|
|
log.info("leads: {} panes recognised {}", leadTerminals.size(), leadTerminals.values());
|
|
}
|
|
// CB-531: on top of the legacy primary.terminal pin, discover leads by the tab labels the
|
|
// operator writes. CB-557 moved the settings onto the lead they describe, so scanning is on
|
|
// whenever a `fleet.leaders:` entry exists — with no leads configured the supplier is a
|
|
// constant and never touches herdr, exactly as a missing `leadScan:` block used to behave.
|
|
// CB-579: each lead now names its own exact `tab:` label, so one scanner discovers every
|
|
// configured lead regardless of how differently their tabs are labelled — the old
|
|
// single-shared-tabPrefix limitation (and its warning) is gone.
|
|
final Supplier<Map<String, String>> leads;
|
|
var leaders = cfg.fleet().leaders();
|
|
if (!leaders.isEmpty()) {
|
|
Set<String> memberSpaces = cfg.profiles().values().stream()
|
|
.map(FleetConfig.Profile::workspace)
|
|
.filter(Objects::nonNull)
|
|
.collect(Collectors.toSet());
|
|
Map<String, String> tabToName = new LinkedHashMap<>();
|
|
leaders.forEach((name, leader) -> {
|
|
if (leader != null && leader.tab() != null && !leader.tab().isBlank()) {
|
|
tabToName.put(leader.tab(), name);
|
|
}
|
|
});
|
|
// One shared rescan cadence: still taken from the first entry, as before — it is an
|
|
// operational cadence, not identity, so there is no correctness reason to give every
|
|
// lead its own scanner.
|
|
int scanIntervalSeconds = leaders.values().iterator().next().scanIntervalSeconds();
|
|
leads = new LeadTabScanner(herdr, tabToName, memberSpaces,
|
|
TimeUnit.SECONDS.toNanos(scanIntervalSeconds), System::nanoTime);
|
|
log.info("lead scan: tabs {} host a lead (rescan every {}s, member spaces {} excluded)",
|
|
tabToName.keySet(), scanIntervalSeconds, memberSpaces);
|
|
} else {
|
|
leads = () -> leadTerminals;
|
|
}
|
|
|
|
// CB-558: start any declared lead that is not already running. After the scanner is built,
|
|
// because both read the same tab labels and the ordering makes that dependency visible; and
|
|
// only when herdr answered, because the launcher's whole safety property is that it can
|
|
// count live leads first — it must never guess and risk a second orchestrator.
|
|
if (herdrUp && !leaders.isEmpty()) {
|
|
int launched = new LeadLauncher(agents, spaces, cfg).ensureLeads();
|
|
if (launched > 0) {
|
|
log.info("lead auto-launch: {} lead(s) started", launched);
|
|
}
|
|
}
|
|
|
|
// CB-548: config-declared architect slots. Config supplies only the stable name → profile
|
|
// map; the terminal → slot binding is owned by the registry and is empty at startup, so no
|
|
// pane resolves to an architect until the later spawn lifecycle binds one. The registry is
|
|
// what CallerResolver resolves against and what that lifecycle will read profiles from;
|
|
// nothing here spawns a slot.
|
|
MemberRegistry members = new MemberRegistry(cfg.fleet());
|
|
sessions.setMemberLifecycle(members);
|
|
if (!members.slots().isEmpty()) {
|
|
log.info("member slots: {} configured {} — none bound yet (a slot is idle until the "
|
|
+ "spawn lifecycle binds a live terminal to it)",
|
|
members.slots().size(), members.slots().keySet());
|
|
}
|
|
|
|
// Status-gated injector (CB-103): the single writer into workers, fed by a poller.
|
|
// The blocking message endpoint (CB-104) is the producer; the poller is inert until then.
|
|
// CB-106: a confirmed turn completion resolves a blocked send whose worker never replied.
|
|
Rendezvous rendezvous = new Rendezvous();
|
|
// CB-578 stage A: classify a completion-fallback scrape that matches a profile's configured
|
|
// usage-limit refusal as BACKEND_EXHAUSTED rather than handing it back as a real answer.
|
|
// Compiled once at startup, keyed by profile name; a profile with no exhaustedPattern is
|
|
// simply absent here, so its workers keep today's completion-fallback behaviour unchanged.
|
|
Map<String, Pattern> exhaustedPatternsByProfile = new LinkedHashMap<>();
|
|
cfg.profiles().forEach((name, profile) -> {
|
|
if (profile.hasExhaustedPattern()) {
|
|
exhaustedPatternsByProfile.put(name, Pattern.compile(profile.exhaustedPattern()));
|
|
}
|
|
});
|
|
ExhaustedPatternLookup exhaustedPatterns = target -> sessions.roster().stream()
|
|
.filter(session -> target.equals(session.terminalId()))
|
|
.findFirst()
|
|
.map(session -> exhaustedPatternsByProfile.get(session.profile()))
|
|
.orElse(null);
|
|
log.info("backend-exhausted classification (CB-578 stage A): {}",
|
|
CompletionResolver.coverage(cfg.profiles().keySet(), exhaustedPatternsByProfile.keySet()));
|
|
// CB-578 stage B: on a classification that actually wins, quarantine the exhausted profile's
|
|
// CREDENTIAL — not the profile name — so a profile sharing that credential (e.g. two models
|
|
// on one OpenAI account) is refused too, not just the one that happened to report it. Reads
|
|
// the profile config live off `config`, so a credentialId edit is hot: no restart needed.
|
|
ExhaustionSink exhaustionSink = (target, reason) -> sessions.roster().stream()
|
|
.filter(session -> target.equals(session.terminalId()))
|
|
.findFirst()
|
|
.map(MemberSession::profile)
|
|
.map(profileName -> config.get().profiles().get(profileName))
|
|
.ifPresent(profile -> {
|
|
String credentialId = profile.effectiveCredentialId();
|
|
quarantine.quarantine(credentialId);
|
|
log.warn("credential '{}' quarantined for {}s (profile '{}' classified "
|
|
+ "BACKEND_EXHAUSTED): {}", credentialId,
|
|
cfg.quarantineCooldownSeconds(), profile.profile(), reason);
|
|
});
|
|
CompletionResolver completion = new CompletionResolver(agents, rendezvous, exhaustedPatterns, exhaustionSink);
|
|
// CB-113: deliver only to an available worker (its MCP is connected), never its boot window.
|
|
// CB-301: the manager's presence bridge records availability and drives SPAWNING → READY.
|
|
MemberPresence presence = sessions.asPresence();
|
|
TurnListener turnListener = new TurnListener() {
|
|
@Override
|
|
public void onTurnComplete(String target) {
|
|
completion.onTurnComplete(target);
|
|
sessions.onTurnComplete(target);
|
|
}
|
|
|
|
@Override
|
|
public boolean hasPostTurnAction(String target) {
|
|
return sessions.hasPostTurnAction(target);
|
|
}
|
|
|
|
@Override
|
|
public boolean onTurnCompleteWithPostAction(String target) {
|
|
completion.resolveBeforePostAction(target);
|
|
return sessions.onTurnCompleteWithPostAction(target);
|
|
}
|
|
|
|
@Override
|
|
public void onDelivered(String target, dev.ltms.fleet.msg.TurnToken token) {
|
|
completion.onDelivered(target, token);
|
|
sessions.onDelivered(target, token);
|
|
}
|
|
|
|
@Override
|
|
public void onTurnFailed(String target) {
|
|
completion.onTurnFailed(target);
|
|
sessions.onTurnFailed(target);
|
|
}
|
|
|
|
@Override
|
|
public void onTurnFailed(String target, String reason) {
|
|
completion.onTurnFailed(target, reason);
|
|
sessions.onTurnFailed(target);
|
|
}
|
|
};
|
|
Injector injector = new Injector(agents, turnListener, deliverableTo(presence, leads),
|
|
presence::forget);
|
|
StatusPoller poller = new StatusPoller(agents, injector, Injector.POLL_INTERVAL_MILLIS);
|
|
poller.start();
|
|
|
|
// 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 = selectReplyInbox(cfg.broker(), System.getenv(), AmqpReplyInbox::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.
|
|
String pinnedPrimaryTerminal = cfg.primary() != null ? cfg.primary().terminal() : null;
|
|
PrimaryRegistry primaryRegistry = new PrimaryRegistry(pinnedPrimaryTerminal);
|
|
// CB-532: `primary.terminal` is superseded and no longer needed for either of its jobs —
|
|
// identity comes from `leaders:`/`leadScan:`, and reply nudges now follow the delegating
|
|
// lead. Say so once at startup rather than leaving a redundant pin to look load-bearing.
|
|
if (pinnedPrimaryTerminal != null && !pinnedPrimaryTerminal.isBlank()) {
|
|
log.warn("primary.terminal is DEPRECATED (CB-532) and can be deleted: identity now comes "
|
|
+ "from leaders:/leadScan:, and reply nudges follow the lead that delegated. "
|
|
+ "It still works, and is still the fallback nudge destination when a restart "
|
|
+ "has lost the delegation map. Its pushReminders/pushBackoffMs stay valid.");
|
|
}
|
|
// CB-307: active push-to-primary loop — nudge the primary when replies land without an
|
|
// open fleet_send. Uses its own lightweight scheduled executor, separate from the injector.
|
|
int maxReminders = cfg.primary() != null ? cfg.primary().remindersOrDefault() : 5;
|
|
long backoffMs = cfg.primary() != null ? cfg.primary().backoffMsOrDefault() : 15_000L;
|
|
var pushScheduler = Executors.newSingleThreadScheduledExecutor(r ->
|
|
Thread.ofVirtual().name("bridge-push-").unstarted(r));
|
|
// 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 = 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
|
|
// an upgraded daemon cannot silently start spending subscription on nudging an idle lead.
|
|
// It has its own single-thread scheduler and holds its own scheduler shutdown via close().
|
|
final LeadHeartbeatLoop heartbeat;
|
|
var heartbeatScheduler = Executors.newSingleThreadScheduledExecutor(r ->
|
|
Thread.ofVirtual().name("bridge-heartbeat-").unstarted(r));
|
|
if (cfg.leadHeartbeat() != null) {
|
|
var hb = cfg.leadHeartbeat();
|
|
heartbeat = new LeadHeartbeatLoop(primaryRegistry, agents, replyInbox, sessions::roster,
|
|
pushLoop, heartbeatScheduler, System::nanoTime,
|
|
TimeUnit.SECONDS.toNanos(hb.idleAfterSeconds()), hb.backoffMs(), hb.quietNudgeCap(),
|
|
metrics);
|
|
heartbeat.start();
|
|
} else {
|
|
heartbeat = null;
|
|
heartbeatScheduler.shutdownNow();
|
|
}
|
|
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
|
|
pushLoop, metrics);
|
|
|
|
// Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller.
|
|
final FleetHealthMonitor healthMonitor;
|
|
var healthScheduler = Executors.newSingleThreadScheduledExecutor(r ->
|
|
Thread.ofVirtual().name("bridge-health-").unstarted(r));
|
|
if (cfg.health() != null && cfg.health().isEnabled()) {
|
|
// CB-580: a member found GONE/NEVER_READY must fail whatever ticket is waiting on it,
|
|
// through the same idempotent target-wide operation CB-516 already uses on release.
|
|
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
|
|
System::nanoTime, cfg.health().intervalOrDefault(), messages::abandon);
|
|
String coverage = FleetHealthMonitor.coverage(true,
|
|
cfg.health().notifications() != null && cfg.health().notifications().configured());
|
|
if ("detection-only".equals(coverage)) {
|
|
log.warn("fleet health: {} (no notification sink configured)", coverage);
|
|
} else {
|
|
log.info("fleet health: {}", coverage);
|
|
}
|
|
healthMonitor.start();
|
|
} else {
|
|
healthMonitor = null;
|
|
healthScheduler.shutdownNow();
|
|
}
|
|
|
|
// CB-520: the reply inbox only consumes for agents this gateway owns. own on acquire,
|
|
// release on teardown. Do this before CB-516 so the inbox is owned before any reply can land.
|
|
sessions.onAcquire(replyInbox::own);
|
|
// CB-516: releasing a worker must fail whatever send was waiting on it. Without this a
|
|
// torn-down delegation kept reporting PENDING until the 30-minute async timeout, and never
|
|
// reached /metrics — the delegation was unresolvable and nothing said so.
|
|
sessions.onRelease(detail -> {
|
|
// CB-578 stage C, acceptance criterion 10: a failed ticket's detail should tell a lead
|
|
// where to re-dispatch onto the same tree, not just that the worker vanished.
|
|
String reason = "the worker session was released before it replied";
|
|
if (detail.worktreePath() != null) {
|
|
reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch()
|
|
+ " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none");
|
|
}
|
|
// CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the
|
|
// member's conversation instead of only re-dispatching a fresh one onto the same files.
|
|
if (detail.agentSessionId() != null) {
|
|
reason += " agentSessionId=" + detail.agentSessionId();
|
|
}
|
|
messages.abandon(detail.terminalId(), reason);
|
|
replyInbox.release(detail.terminalId());
|
|
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
|
|
});
|
|
|
|
// MCP server face (CB-105): fleet_send/fleet_reply/fleet_status, mounted at /mcp.
|
|
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
|
|
ConnectionIdentity identity = new ConnectionIdentity(
|
|
new PaneLocator(herdr), new LsofPeerPidLookup(), new LsofProcessCwdLookup());
|
|
|
|
// CB-501: one resolver behind both entry paths. Worker identity still comes from the
|
|
// connection and is never token-gated, so enabling token mode cannot lock the fleet out.
|
|
final CallerResolver callers;
|
|
if (cfg.auth().tokenMode()) {
|
|
String token = System.getenv(cfg.auth().tokenEnv());
|
|
if (token == null || token.isBlank()) {
|
|
throw new IllegalStateException("auth.mode=token but env var " + cfg.auth().tokenEnv()
|
|
+ " is unset or empty — export it before starting bridged");
|
|
}
|
|
callers = CallerResolver.withLeadsAndMembers(identity, true, token, leads, members);
|
|
log.info("auth: token mode (bearer required for non-worker callers, env {})",
|
|
cfg.auth().tokenEnv());
|
|
} else {
|
|
callers = CallerResolver.withLeadsAndMembers(identity, false, null, leads, members);
|
|
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
|
|
}
|
|
|
|
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 FleetMcp.HealthCoverageSource(() -> {
|
|
var health = config.get().health();
|
|
return FleetHealthMonitor.coverage(health != null && health.isEnabled(),
|
|
health != null && health.notifications() != null && health.notifications().configured());
|
|
}),
|
|
new FleetMcp.QuarantineSource(profile -> {
|
|
var configured = config.get().profiles().get(profile);
|
|
return configured == null ? null : configured.effectiveCredentialId();
|
|
}, quarantine));
|
|
|
|
// 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.
|
|
final ConfigWatcher configWatcher;
|
|
if (cfg.configReload() != null && cfg.configReload().isEnabled()) {
|
|
configWatcher = new ConfigWatcher(config, cfg.configReload().intervalSeconds());
|
|
configWatcher.start();
|
|
} else {
|
|
configWatcher = null;
|
|
}
|
|
|
|
// CB-303 part 3: single ordered shutdown hook. Drain sessions first while herdr is still
|
|
// open (so releases reach the daemon), then stop poller/message/mcp/reaper, and close herdr
|
|
// last. This replaces the earlier independent hooks that could race and close herdr early.
|
|
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
|
|
sessions.close(cfg.lifecycle() != null ? cfg.lifecycle().drainTimeoutSeconds() : null);
|
|
poller.stop();
|
|
messages.close();
|
|
pushLoop.close();
|
|
if (heartbeat != null) heartbeat.close(); // CB-551: stop the idle-lead heartbeat scheduler
|
|
if (healthMonitor != null) healthMonitor.stop();
|
|
if (configWatcher != null) configWatcher.stop(); // CB-559: stop polling the config file
|
|
mcp.close();
|
|
if (reaper != null) reaper.stop();
|
|
// Release the broker connection last among message resources (no-op for the in-memory inbox).
|
|
if (replyInbox instanceof AutoCloseable closeable) {
|
|
try {
|
|
closeable.close();
|
|
} catch (Exception e) {
|
|
log.debug("reply inbox close: {}", e.toString());
|
|
}
|
|
}
|
|
herdr.close();
|
|
}));
|
|
|
|
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 {}",
|
|
cfg.bind().host(), cfg.bind().port(), socket);
|
|
}
|
|
|
|
/**
|
|
* The {@link Injector}'s readiness gate (CB-534): a target is deliverable if it is a spawned
|
|
* member whose agent has connected the bridge MCP, <em>or</em> a lead.
|
|
*
|
|
* <p>The gate exists for one reason — to hold a delivery out of a <em>spawned</em> member's boot
|
|
* window, where herdr already reports {@code idle} but the TUI would drop an injected paste. That
|
|
* 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 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
|
|
* lead→lead send sat on the gate for {@code READINESS_GRACE_POLLS} (~60s) and then failed
|
|
* having never been typed into the pane.
|
|
*
|
|
* <p>The lead set is read through the supplier on each call rather than snapshotted, so a lead
|
|
* discovered by {@code leadScan} after startup becomes deliverable without a restart.
|
|
*/
|
|
static Predicate<String> deliverableTo(MemberPresence presence, Supplier<Map<String, String>> leads) {
|
|
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);
|
|
}
|
|
|
|
/**
|
|
* 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 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(FleetConfig)} is the logging caller.
|
|
*/
|
|
static Map<String, List<String>> requiredSecretEnvVars(FleetConfig cfg) {
|
|
Map<String, List<String>> requiredBy = new LinkedHashMap<>();
|
|
cfg.profiles().forEach((name, profile) -> {
|
|
if (!profile.isSubscription()) {
|
|
requiredBy.computeIfAbsent(profile.tokenEnv(), _ -> new ArrayList<>())
|
|
.add("profile '" + name + "' tokenEnv");
|
|
}
|
|
if (profile.hasGitToken()) {
|
|
requiredBy.computeIfAbsent(profile.gitTokenEnv(), _ -> new ArrayList<>())
|
|
.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;
|
|
}
|
|
|
|
/**
|
|
* CB-594: log, by name only, which required env vars (see {@link #requiredSecretEnvVars}) are
|
|
* set in the daemon's own process environment — the environment every profile's {@code
|
|
* tokenEnv}/{@code gitTokenEnv} is read from at spawn time (see
|
|
* {@code HerdrPeerLauncher.resolveEnv}). Never logs a value, a prefix, or a length.
|
|
*
|
|
* <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(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");
|
|
return;
|
|
}
|
|
Map<String, String> env = System.getenv();
|
|
requiredBy.forEach((varName, sources) -> {
|
|
String value = env.get(varName);
|
|
if (value != null && !value.isBlank()) {
|
|
log.info("startup secret {}: set ({})", varName, String.join(", ", sources));
|
|
} else {
|
|
log.warn("startup secret {}: MISSING ({}) — the daemon will start anyway, and this "
|
|
+ "failure stays invisible until a worker actually needs it. Fix "
|
|
+ "${SHARED_ENV}/tools/secrets.sh and restart bridged from a LOGIN "
|
|
+ "shell (see scripts/redeploy-bridged.sh).",
|
|
varName, String.join(", ", sources));
|
|
}
|
|
});
|
|
}
|
|
|
|
/**
|
|
* CB-596: {@code known:} empty (block absent entirely, or present but empty) means {@link
|
|
* 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
|
|
* start over it — see {@link #reportRequiredSecrets} for why a daemon that boots and says
|
|
* what is wrong beats one that will not boot at all.
|
|
*
|
|
* <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(FleetConfig cfg) {
|
|
FleetConfig.MemberCredentials creds = cfg.memberCredentials();
|
|
if (creds != null && !creds.known().isEmpty()) {
|
|
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 fleetd.example.yaml — and restart.");
|
|
}
|
|
|
|
/**
|
|
* Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504).
|
|
*
|
|
* @return true if herdr answered, false if it never did
|
|
*/
|
|
private static boolean awaitHerdr(HerdrClient herdr) {
|
|
long deadline = System.nanoTime() + HERDR_WAIT_SECONDS * 1_000_000_000L;
|
|
boolean waited = false;
|
|
while (true) {
|
|
try {
|
|
herdr.call("ping");
|
|
if (waited) {
|
|
log.info("herdr is up");
|
|
}
|
|
return true;
|
|
} catch (HerdrException e) {
|
|
if (System.nanoTime() >= deadline) {
|
|
return false;
|
|
}
|
|
if (!waited) {
|
|
log.info("waiting up to {}s for the herdr socket…", HERDR_WAIT_SECONDS);
|
|
waited = true;
|
|
}
|
|
try {
|
|
Thread.sleep(HERDR_WAIT_POLL_MILLIS);
|
|
} catch (InterruptedException ie) {
|
|
Thread.currentThread().interrupt();
|
|
return false;
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
private Fleetd() {
|
|
}
|
|
}
|