CB-185: route message status by target
This commit is contained in:
@@ -446,7 +446,7 @@ public final class Fleetd {
|
||||
heartbeat = null;
|
||||
heartbeatScheduler.shutdownNow();
|
||||
}
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
|
||||
MessageService messages = new MessageService(router, injector, rendezvous, replyInbox,
|
||||
pushLoop, metrics);
|
||||
|
||||
// Health is a slow whole-fleet observer. Keep it separate from the 250ms delivery poller.
|
||||
|
||||
@@ -2,6 +2,7 @@ package dev.ltms.fleet.msg;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.AgentStatus;
|
||||
import dev.ltms.fleet.herdr.HerdrRouter;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.metrics.FleetMetrics;
|
||||
import dev.ltms.fleet.metrics.Metrics;
|
||||
@@ -189,6 +190,7 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
private final AgentControl agents;
|
||||
private final HerdrRouter router;
|
||||
private final Injector injector;
|
||||
private final Rendezvous rendezvous;
|
||||
private final ReplyInbox inbox;
|
||||
@@ -257,6 +259,7 @@ public final class MessageService {
|
||||
MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox,
|
||||
ReplyPushLoop pushLoop, Metrics metrics, LongSupplier nowNanos) {
|
||||
this.agents = agents;
|
||||
this.router = null;
|
||||
this.injector = injector;
|
||||
this.rendezvous = rendezvous;
|
||||
this.inbox = inbox;
|
||||
@@ -265,6 +268,20 @@ public final class MessageService {
|
||||
this.nowNanos = nowNanos;
|
||||
}
|
||||
|
||||
public MessageService(HerdrRouter router, Injector injector, Rendezvous rendezvous, ReplyInbox inbox,
|
||||
ReplyPushLoop pushLoop, Metrics metrics) {
|
||||
this.agents = null;
|
||||
this.router = router;
|
||||
this.injector = injector;
|
||||
this.rendezvous = rendezvous;
|
||||
this.inbox = inbox;
|
||||
this.pushLoop = pushLoop;
|
||||
this.metrics = metrics;
|
||||
this.nowNanos = System::nanoTime;
|
||||
}
|
||||
|
||||
private AgentControl agentsFor(String target) { return router != null ? router.agentsFor(target) : agents; }
|
||||
|
||||
/** Create with an explicit {@link ReplyInbox} and no push loop. */
|
||||
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous, ReplyInbox inbox) {
|
||||
this(agents, injector, rendezvous, inbox, null);
|
||||
@@ -277,7 +294,7 @@ public final class MessageService {
|
||||
|
||||
/** Current lifecycle status of a worker (the {@code GET /sessions/{id}/status} surface). */
|
||||
public AgentStatus status(String target) {
|
||||
return agents.status(target);
|
||||
return agentsFor(target).status(target);
|
||||
}
|
||||
|
||||
/** Read-only delegation fact for fleet views. */
|
||||
@@ -788,7 +805,7 @@ public final class MessageService {
|
||||
/** Best-effort live worker status for a pending poll; never throws (a lookup error is just noise). */
|
||||
private String liveStatus(String target) {
|
||||
try {
|
||||
return agents.status(target).name().toLowerCase();
|
||||
return agentsFor(target).status(target).name().toLowerCase();
|
||||
} catch (RuntimeException e) {
|
||||
return "unknown";
|
||||
}
|
||||
|
||||
@@ -0,0 +1,17 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
|
||||
class FleetdHerdrControlConstructionTest {
|
||||
@Test
|
||||
void fleetdDelegatesStatefulControlsToTheRouter() throws Exception {
|
||||
// AgentControl caches paneByTerminal, so the router must be its only production factory.
|
||||
String source = Files.readString(Path.of("src/main/java/dev/ltms/fleet/Fleetd.java"));
|
||||
assertFalse(source.contains("new AgentControl("));
|
||||
assertFalse(source.contains("new WorkspaceControl("));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user