9b50dd69d8
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.
94 lines
3.6 KiB
Java
94 lines
3.6 KiB
Java
package dev.ltms.fleet.inject;
|
|
|
|
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;
|
|
|
|
import java.util.Set;
|
|
|
|
/**
|
|
* Drives the {@link Injector} by sampling each active worker's {@code agent_status} and
|
|
* feeding it in. A single virtual-thread loop polls only targets that have work outstanding,
|
|
* so an idle bridge does no herdr traffic.
|
|
*
|
|
* <p>This is the Stage-1 gate signal. It can later be replaced (or fronted) by a herdr
|
|
* {@code events.subscribe} stream without touching the {@link Injector} — the injector only
|
|
* consumes {@code onStatus} calls, however they are produced.
|
|
*/
|
|
public final class StatusPoller {
|
|
|
|
private static final Logger log = LoggerFactory.getLogger(StatusPoller.class);
|
|
|
|
private final AgentControl agents;
|
|
private final Injector injector;
|
|
private final StatusRefiner refiner;
|
|
private final long intervalMillis;
|
|
private volatile boolean running;
|
|
private Thread thread;
|
|
|
|
public StatusPoller(AgentControl agents, Injector injector, long intervalMillis) {
|
|
this(agents, injector, new StatusRefiner(agents), intervalMillis);
|
|
}
|
|
|
|
public StatusPoller(AgentControl agents, Injector injector, StatusRefiner refiner,
|
|
long intervalMillis) {
|
|
this.agents = agents;
|
|
this.injector = injector;
|
|
this.refiner = refiner;
|
|
this.intervalMillis = intervalMillis;
|
|
}
|
|
|
|
/** Start the polling loop on a virtual thread. Idempotent. */
|
|
public synchronized void start() {
|
|
if (running) return;
|
|
running = true;
|
|
thread = Thread.ofVirtual().name("status-poller").start(this::loop);
|
|
log.info("status poller started (interval {}ms)", intervalMillis);
|
|
}
|
|
|
|
private void loop() {
|
|
while (running) {
|
|
Set<String> active = injector.activeTargets();
|
|
for (String target : active) {
|
|
if (!running) return;
|
|
try {
|
|
// herdr's agent_status can misreport a settled worker as `unknown`; refine it
|
|
// against the pane content before it drives delivery/completion (CB-115).
|
|
AgentStatus status = refiner.refine(target, agents.status(target));
|
|
injector.onStatus(target, status);
|
|
} catch (HerdrException e) {
|
|
// The worker's agent is gone — stop trying and unblock its waiters.
|
|
if (e.code() != null && e.code().endsWith("_not_found")) {
|
|
log.debug("target {} gone; dropping its queue", target);
|
|
injector.drop(target, e);
|
|
} else {
|
|
log.debug("status poll for {} failed (will retry): {}", target, e.getMessage());
|
|
}
|
|
} catch (RuntimeException e) {
|
|
// Never let one target's unexpected error (e.g. an odd agent.get shape) kill
|
|
// the single poller thread and stall injection for every worker.
|
|
log.warn("unexpected error polling {}; skipping this round", target, e);
|
|
}
|
|
}
|
|
sleep();
|
|
}
|
|
}
|
|
|
|
private void sleep() {
|
|
try {
|
|
Thread.sleep(intervalMillis);
|
|
} catch (InterruptedException e) {
|
|
Thread.currentThread().interrupt();
|
|
running = false;
|
|
}
|
|
}
|
|
|
|
/** Stop the polling loop. Idempotent. */
|
|
public synchronized void stop() {
|
|
running = false;
|
|
if (thread != null) thread.interrupt();
|
|
}
|
|
}
|