Compare commits

..

9 Commits

Author SHA1 Message Date
Dai Ha a42b12440c Merge PR #649: fleetd #612 Shape A r9+r11 — pin capacitySource, healthCoverageSource, coordinator peers via a real fleet_list round-trip
CI / shell-tests (push) Failing after 9s
CI / contract (push) Successful in 58s
CI / build (push) Failing after 1m52s
Test-only, 243 lines, one new file. ZERO production change — this is the point
of the PR, and it replaces PR #647, which added three public accessors to
FleetMcp on a false premise.

#647 argued a same-JVM test caller can never reach fleet_list's READ gate,
because FleetdAssembly hardcodes a real LsofPeerPidLookup, so the caller
resolves ANONYMOUS. The premise about lsof is true; the conclusion is not. In
CallerResolver.resolve, once c.terminal() == null the next branch is:

    if (tokenMode) {
        return presentedTokenMatches(authorizationHeader)
                ? Principal.primary(c.pid()) : Principal.anonymous();
    }

c.resolved() and c.scanComplete() guard only the LATER loopback-trust branch,
which token mode returns before reaching. So under auth.mode: token a bearer
token resolves to PRIMARY with no pid lookup on the path. This PR does exactly
that: real McpSyncClient callTool("fleet_list") over a real transport with an
Authorization: Bearer header, a faked leadMailboxOpener so no broker is touched,
and a real coordinator: block with a non-empty peers list. No reflection.

Lead verification, measured myself (not taken from the worker's report):

  git diff --stat origin/main -- fleetd/src/main/java/   -> EMPTY
  mutate :481 Fleetd.capacitySource(...) -> CapacitySource.none()
      -> Tests run: 3, Failures: 1
         fleetListReportsTheAssembledCapacitySource:222
         the other two tests stayed GREEN

The failure message carries the live fleet_list JSON body, showing
healthCoverage, loopHealth and coordinator.peers present with capacity absent —
so the round-trip, the PRIMARY resolution and the coordinator visibility are all
real, and the mutation removed exactly one thing.

Worker also reported, each with a grep -c anchor of 1 restored: health mutation
-> 1 failure named fleetListReportsTheAssembledHealthCoverageSource; peers
mutation -> 1 failure named fleetListReportsTheAssembledCoordinatorPeers; final
mvn clean install Tests run: 1899, Failures: 0 / BUILD SUCCESS. I reproduced the
capacity cycle only; the other two are the worker's measurement, not mine.

Carries forward #647's genuine find: the peers ternary at :489 is
inert-equals-absent, so the test configures a real coordinator: block with peers.

Detail on #612 and #647.
2026-10-02 04:53:56 +02:00
Dai Ha a06426c33c Merge PR #648: fleetd #612 Shape A r4 — pin quarantineSource + outageSource through BOTH operator windows
Test-only, 368 lines, one new file. No production change.

Lead verification, measured myself in a throwaway detached worktree (not taken
from the worker's report):

  starve FleetApp pass site :529   -> Tests run: 4, Failures: 2
                                      (...ByTheRealAssembledFleetApp:333, :363)
                                      both ...FleetMcp tests stayed GREEN
  starve FleetMcp pass sites :484/:486 -> Tests run: 4, Failures: 2
                                      (...ByTheRealAssembledFleetMcp:319, :349)
                                      both ...FleetApp tests stayed GREEN

So the two windows are pinned independently. Failure messages carry the live
JSON response body from a real round-trip, not a source-text match.

Why this is worth merging — the REST window was completely unpinned. With this
PR's test parked and :529 starved, the FULL suite reported:

  Tests run: 1892, Failures: 0, Errors: 0, Skipped: 0 / BUILD SUCCESS

Zero pre-existing tests notice the REST window losing its sources. An operator
reads GET /profiles exactly when the MCP mount is down. Arithmetic control:
1892 + this PR's 4 = 1896 = main at 6539efe.

Known caveat, recorded not fixed: the MCP-side quarantine assertion (:313) is
DUPLICATE coverage. A reviewer measured on clean origin/main that starving the
MCP quarantine arg already fails three pre-existing tests
(FleetdBackendQuarantineAssemblyTest:167, FleetdExhaustedPatternAssemblyTest:202,
FleetdOpenCodeExhaustionForwardingAssemblyTest:182), so that site was already
pinned and the PR's claim otherwise is wrong. Low severity, left in place as a
valid fleet_profiles output assertion. Detail on #648 and #612.

Reviewed by two reviewers against the diff (soundness: no issue; coupling: the
duplicate-coverage finding above).
2026-10-02 04:53:39 +02:00
Dai Ha 28ea0de575 fleetd #612: cover fleet list reporting sources
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 1m9s
CI / build (pull_request) Failing after 2m10s
2026-10-02 04:42:07 +02:00
Dai Ha 91792e11fc fleetd #612 Shape A unit r4: pin quarantineSource + outageSource through both operator windows
CI / shell-tests (pull_request) Failing after 9s
CI / contract (pull_request) Successful in 1m29s
CI / build (pull_request) Failing after 1m56s
FleetdAssembly.java's quarantineSource (:471-472) and outageSource (:473-476)
each feed two consumers: FleetMcp (fleet_profiles, :484/:486) and FleetApp
(GET /profiles, :529). No existing test distinguished the two windows for
either source.

New test drives the real FleetdAssembly.assembleAndStart, classifies a real
exhaustion/outage through the real CompletionResolver, and reads the result
back through a real McpSyncClient (fleet_profiles) and a real HttpClient
(GET /profiles), both authenticated via token-mode auth (sidesteps the
in-JVM pid-resolution dead end). No source text is read; no production code
changed.

Verified with six mutation cycles (3 per site x 2 sites: FleetMcp starved,
FleetApp starved, mis-wire with a disconnected collaborator), each run
against the full unfiltered suite and reverted after confirming the
expected test(s) alone went red. The quarantineSource FleetMcp-starve and
mis-wire cycles also trip three pre-existing tests that read the live
BackendQuarantine via FleetMcp#quarantineSource() for their own unrelated
assertions - a pre-existing incidental coupling, not newly introduced here.

Out of scope, noted per the ticket's dispatch comment: loopHealthSource
(FleetdAssembly.java:478) shares this same two-consumer shape and is
already assigned to a separate unit, r10.
2026-10-02 04:20:53 +02:00
ltms 6539efe9fa Merge PR #645: fleetd #612 Shape A r10 — pin loopHealthSource at BOTH consumers
CI / shell-tests (push) Failing after 7s
CI / contract (push) Successful in 1m34s
CI / build (push) Failing after 2m11s
Pins FleetdAssembly's loopHealth local at both of its pass sites: FleetMcp (:483, the fleet_list source) and FleetApp (:529, the real /healthz body). Two independent assertions, so starving one site leaves the other green.

Lead verification, independent of the implementer's own proof, in a throwaway detached worktree at 68397f7:
- :483 starved only -> Tests run: 2, Failures: 1. RED: fleetListUsesTheRunningLoopsInTheRealAssembledMcp, expected: <RUNNING> but was: <STOPPED>. REST test stayed GREEN.
- :529 starved only -> Tests run: 2, Failures: 1. RED: healthzUsesTheRunningLoopsInTheRealAssembledApp, with the real body {"loopHealth":{"sessionReaper":"STOPPED","statusPoller":"STOPPED"}}. MCP test stayed GREEN.
- full suite with the PR: Tests run: 1896, Failures: 0, Errors: 0 — BUILD SUCCESS (1894 baseline + 2).

The REST half binds port 0 (ephemeral) via runtime.app().start("127.0.0.1", 0), not the configured 8765, so it cannot clash with the live daemon. Teardown stops the bound app, closes the runtime, and asserts the shutdown hook was registered and the herdr client actually closed.

I chased the implementer's honestly-reported anomaly (an unrelated ClaudeCodeLauncherTest.noFixtureSeededTheDefaultClaudeJson failure in one cycle, and a suite total that moved between cycles). It is NOT caused by this PR: a baseline run at 68397f7 with this PR absent added the same 3 temp-dir project entries to the real ~/.claude.json as the run with it applied (65->68 without, 68->71 with), and ClaudeCodeLauncherTest passed in both. Filed separately.

Test-only diff, no production code touched.
2026-10-02 04:15:04 +02:00
ltms 68397f78d5 Merge PR #644: fleetd #612 Shape A r12 — pin the assembled turn registrar
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 1m32s
CI / build (push) Failing after 1m57s
Pins FleetdAssembly.java:349 behaviourally. The test installs a throwing wrapper on the real assembled Injector's listener to construct the narrow between-delivery-and-completion window, then asserts runtime.completion() still resolves the registered waiter. No sleep: it advances a controllable nano clock. Teardown runs the captured shutdown hook and asserts herdr actually closed.

Lead verification, independent of the implementer's own proof, in a throwaway detached worktree at dac5f88:
- unmutated: Tests run: 1, Failures: 0 — BUILD SUCCESS
- :349 registrar -> (_, _) -> { }: Failures: 1 — "must wire the Injector registrar to this runtime's real CompletionResolver" expected: <true> but was: <false>
- mis-wire I built myself (differs from the implementer's): Fleetd.turnRegistrar(new CompletionResolver(agents, new Rendezvous(), ...)) — a fresh Rendezvous so the registry genuinely differs: Failures: 1, same assertion
- reverted, full suite: Tests run: 1894, Failures: 0, Errors: 0 — BUILD SUCCESS, 58s (1892 baseline + r5 + r12)

Test-only diff, no production code touched. Reflection is used only to install the throwing listener; that is how the failure window is constructed, not how the assertion is made. No reviewer fan-out: member capacity is committed to the remaining Shape A implementers.
2026-10-02 04:09:29 +02:00
Dai Ha af901ff1d2 fleetd #612: pin assembled loop health
CI / shell-tests (pull_request) Failing after 10s
CI / contract (pull_request) Successful in 1m1s
CI / build (pull_request) Failing after 2m23s
2026-10-02 04:08:02 +02:00
ltms dac5f88812 Merge PR #643: fleetd #612 Shape A r5 — pin FleetdAssembly's leadConfigDirSource call site
CI / shell-tests (push) Failing after 8s
CI / contract (push) Successful in 55s
CI / build (push) Failing after 2m14s
Pins the #602/#606 call site behaviourally: drives the real FleetdAssembly.assembleAndStart and asserts the assembled LeadConfigDirSource resolves a real configured configDir, which none() cannot produce.

Lead verification, run independently of the implementer's own proof, in a throwaway detached worktree at 141ae3b:
- unmutated: Tests run: 1, Failures: 0 — BUILD SUCCESS
- FleetdAssembly.java:488 -> FleetMcp.LeadConfigDirSource.none(): Tests run: 1, Failures: 1 — expected: </mnt/fake-lead-configdir> but was: <null>
- reverted, full suite: Tests run: 1893, Failures: 0, Errors: 0 — BUILD SUCCESS, 59s (baseline 1892)

Test-only diff, no production code touched. No reviewer fan-out was run: all member capacity is committed to the five Shape A implementers.
2026-10-02 04:05:50 +02:00
Dai Ha 2e663e5968 fleetd #612: pin assembled turn registrar
CI / shell-tests (pull_request) Failing after 6s
CI / contract (pull_request) Successful in 1m28s
CI / build (pull_request) Failing after 1m51s
2026-10-02 04:04:54 +02:00
4 changed files with 975 additions and 0 deletions
@@ -0,0 +1,180 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.inject.LoopWatchdog;
import dev.ltms.fleet.mcp.FleetMcp;
import io.javalin.Javalin;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.lang.reflect.Field;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #612 Shape A, rank 10: {@link FleetdAssembly} creates one {@link FleetMcp.LoopHealthSource}
* from the real started {@code StatusPoller} and {@code SessionReaper}, then gives it to two operator
* windows. These tests reach the real assembled objects through {@link FleetdRuntime}, rather than
* building a second source beside them. A hardcoded {@code RUNNING} source would pass a simple
* "running" test, so the mutation proof also mis-wires the source to never-started loops: both
* windows must then report {@code STOPPED} and these assertions go red.
*/
class FleetdAssemblyLoopHealthTest {
private static final class RecordingResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
Runnable shutdownHook;
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> new dev.ltms.fleet.msg.InMemoryReplyInbox();
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException("no coordinator is configured");
};
}
@Override
public Runnable herdrPollWait() {
return () -> {
throw new UnsupportedOperationException("healthy FakeHerdr must not be polled");
};
}
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
}
@Override
public LongSupplier wallClockNanos() {
return System::nanoTime;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
// Bind runtime.app() to an ephemeral port only in the REST assertion below.
}
}
private final HttpClient http = HttpClient.newHttpClient();
private RecordingResourcePorts ports;
private FleetdRuntime runtime;
private Javalin boundApp;
@AfterEach
void tearDown() {
if (boundApp != null) {
boundApp.stop();
}
if (runtime != null) {
runtime.close();
}
if (ports != null) {
assertNotNull(ports.shutdownHook, "the assembly must register its shutdown hook");
assertTrue(ports.herdr.closed, "FleetdRuntime.close must close the real assembled herdr client");
}
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, """
bind:
host: 127.0.0.1
port: 8765
lifecycle:
idleTtlSeconds: 600
idleSleepGuard:
enabled: false
health:
enabled: false
broker:
uri: "amqp://fake-test-broker/vh"
""");
return FleetConfig.load(file);
}
private void assemble(Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir);
ports = new RecordingResourcePorts();
runtime = FleetdAssembly.assembleAndStart(
new AssemblyInputs(cfg, new ConfigRef(dir.resolve("fleetd.yaml"), cfg),
new SubscriptionGuard(cfg.guard().hostSet())), ports);
}
@Test
void fleetListUsesTheRunningLoopsInTheRealAssembledMcp(@TempDir Path dir) throws Exception {
assemble(dir);
// The field is the exact source captured by FleetMcp's fleet_list handler. Reflection is
// necessary because FleetMcp has no public source accessor; it is not a source-text check.
FleetMcp.LoopHealthSource loopHealth = loopHealthOf(runtime.mcp());
assertRunning(loopHealth, "FleetMcp's real fleet_list source");
}
@Test
void healthzUsesTheRunningLoopsInTheRealAssembledApp(@TempDir Path dir) throws Exception {
assemble(dir);
boundApp = runtime.app().start("127.0.0.1", 0);
HttpRequest request = HttpRequest.newBuilder(
URI.create("http://127.0.0.1:" + boundApp.port() + "/healthz")).GET().build();
HttpResponse<String> response = http.send(request, HttpResponse.BodyHandlers.ofString());
assertEquals(200, response.statusCode(), response.body());
assertTrue(response.body().contains("\"statusPoller\":\"RUNNING\""), response.body());
assertTrue(response.body().contains("\"sessionReaper\":\"RUNNING\""), response.body());
}
private static FleetMcp.LoopHealthSource loopHealthOf(FleetMcp mcp) throws Exception {
Field field = FleetMcp.class.getDeclaredField("loopHealth");
field.setAccessible(true);
return (FleetMcp.LoopHealthSource) field.get(mcp);
}
private static void assertRunning(FleetMcp.LoopHealthSource source, String consumer) {
assertEquals(LoopWatchdog.State.RUNNING, source.statusPoller().get(),
consumer + " must report the started real StatusPoller as RUNNING");
assertEquals(LoopWatchdog.State.RUNNING, source.sessionReaper().get(),
consumer + " must report the started real SessionReaper as RUNNING");
}
}
@@ -0,0 +1,184 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.AgentStatus;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.inject.Injector;
import dev.ltms.fleet.inject.TurnListener;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.TurnToken;
import io.javalin.Javalin;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import java.lang.reflect.Field;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #612 rank 12 — {@code FleetdAssembly} supplies the {@link Injector}'s {@code TurnRegistrar}
* with {@code Fleetd.turnRegistrar(completion)}. The normal delivery path cannot distinguish that
* registrar from {@code TurnRegistrar.NOOP}: {@link dev.ltms.fleet.inject.CompletionResolver#onDelivered}
* registers the same turn shortly afterwards. The distinction matters when a delivery listener throws
* after the pane received the message but before normal completion runs. The registrar must already have
* registered the waiter with the REAL assembled resolver, so a later completion can still resolve it.
*
* <p>The test installs a throwing wrapper around the real assembled listener. It does not replace the
* registrar or resolver. This constructs the narrow failure condition without a sleep, then observes the
* real resolver through {@link FleetdRuntime#completion()}. An inert registrar, or a registrar wired to a
* throwaway resolver, leaves this waiter's turn absent from the real resolver and makes the final assertion
* fail.
*/
class FleetdAssemblyTurnRegistrarBehaviouralTest {
private static final String TARGET = "term_a";
private static final class ControllableResourcePorts implements ResourcePorts {
final FakeHerdr herdr = new FakeHerdr();
final AtomicLong nowNanos = new AtomicLong(1_000_000_000L);
Runnable shutdownHook;
void advanceSeconds(long seconds) {
nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(seconds));
}
@Override
public Map<String, String> environment() {
return Map.of();
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> {
throw new UnsupportedOperationException("no broker: block is configured");
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException("no coordinator: block is configured");
};
}
@Override
public LongSupplier nanoClock() {
return nowNanos::get;
}
@Override
public LongSupplier wallClockNanos() {
return nowNanos::get;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
// Do not bind a real port in this assembly test.
}
@Override
public Runnable herdrPollWait() {
return () -> {
throw new UnsupportedOperationException("FakeHerdr is healthy; no poll wait is expected");
};
}
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, """
bind:
host: 127.0.0.1
port: 8765
idleSleepGuard:
enabled: false
fleet:
leaders:
primary:
tab: "lead: primary"
profile: sonnet
profiles:
sonnet:
subscription: true
argv: ["ccs", "sonnet"]
""");
return FleetConfig.load(file);
}
@Test
void realAssembledResolverStillResolvesAfterDeliveredListenerThrows(@TempDir Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir);
ControllableResourcePorts ports = new ControllableResourcePorts();
ports.herdr.withTab("w2", "w2:t7", "lead: primary");
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg,
new ConfigRef(dir.resolve("fleetd.yaml"), cfg), new SubscriptionGuard(cfg.guard().hostSet())), ports);
try {
Injector injector = runtime.injector();
installThrowingDeliveredListener(injector);
CompletableFuture<Rendezvous.Resolution> waiter = new CompletableFuture<>();
injector.enqueue(TARGET, "brief", new TurnToken(TARGET, waiter));
IllegalStateException thrown = assertThrows(IllegalStateException.class,
() -> injector.onStatus(TARGET, AgentStatus.IDLE),
"control: delivery must reach the installed listener and it must throw after delivery");
assertTrue(thrown.getMessage().contains("listener failure"), thrown::getMessage);
ports.herdr.readText("worker report after the listener failure");
ports.advanceSeconds(3);
runtime.completion().resolveBeforePostAction(TARGET);
assertTrue(waiter.isDone(),
"FleetdAssembly must wire the Injector registrar to this runtime's real CompletionResolver: "
+ "after a delivered listener throws, resolveBeforePostAction must still find and "
+ "resolve the registered waiter");
} finally {
assertTrue(ports.shutdownHook != null, "control: assembly must capture its shutdown hook");
ports.shutdownHook.run();
assertTrue(ports.herdr.closed, "teardown control: the captured shutdown hook must close herdr");
}
}
private static void installThrowingDeliveredListener(Injector injector) throws Exception {
Field field = Injector.class.getDeclaredField("turnListener");
field.setAccessible(true);
field.set(injector, new TurnListener() {
@Override
public void onTurnComplete(String target) {
}
@Override
public void onDelivered(String target, TurnToken token) {
throw new IllegalStateException("listener failure after delivery");
}
});
}
}
@@ -0,0 +1,243 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.msg.LeadChannelHandle;
import dev.ltms.fleet.msg.LeadMessage;
import dev.ltms.fleet.msg.ReplyInbox;
import io.javalin.Javalin;
import io.modelcontextprotocol.client.McpClient;
import io.modelcontextprotocol.client.McpSyncClient;
import io.modelcontextprotocol.client.transport.HttpClientStreamableHttpTransport;
import io.modelcontextprotocol.spec.McpClientTransport;
import io.modelcontextprotocol.spec.McpSchema;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
import org.junit.jupiter.api.io.TempDir;
import java.net.http.HttpRequest;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.List;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #612 Shape A ranks 9 and 11: drives the real {@code fleet_list} MCP route through a
* token-authenticated primary caller. The assertions read the response from the exact {@link
* dev.ltms.fleet.mcp.FleetMcp} instance assembled by {@link FleetdAssembly}, rather than a source
* scrape or a separately built reporting source.
*
* <p>The coordinator fixture contains a peer on purpose. Both a missing coordinator and an
* incorrectly wired peers argument can render as an empty list, so the non-empty peer assertion
* distinguishes the real call-site value from that inert result. Its fake mailbox never contacts a
* broker.
*/
class FleetdListReportingSourcesAssemblyTest {
private static final String TOKEN = "fleetd-r9-r11-test-token";
private static final String TOKEN_ENV = "FLEETD_R9_TEST_TOKEN";
private static final String PROFILE = "capacity-profile";
private static final String PEER = "peer-fleet";
private static final class FakeLeadChannel implements LeadChannelHandle {
@Override
public void publish(String toCoordId, LeadMessage message) {
}
@Override
public List<LeadMessage> peek() {
return List.of();
}
@Override
public void ack(String msgId) {
}
@Override
public String selfCoordId() {
return "test-lead";
}
@Override
public boolean heldDurable() {
return true;
}
@Override
public MailboxState inspect(String coordId) {
return MailboxState.unknown(coordId);
}
@Override
public void close() {
}
}
private static final class TestResourcePorts implements ResourcePorts {
private final FakeHerdr herdr = new FakeHerdr();
private Runnable shutdownHook;
@Override
public Map<String, String> environment() {
return Map.of(TOKEN_ENV, TOKEN);
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
return (uri, prefetch) -> new ReplyInbox() {
@Override public void own(String target) { }
@Override public void release(String target) { }
@Override public void publish(String target, String msgId, String content) { }
@Override public List<InboxMessage> peek(String target) { return List.of(); }
@Override public boolean ack(String target, String msgId) { return false; }
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
return (uri, selfCoordId, prefetch) -> new FakeLeadChannel();
}
@Override
public LongSupplier nanoClock() {
return System::nanoTime;
}
@Override
public LongSupplier wallClockNanos() {
return System::nanoTime;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
// The test binds runtime.app() itself, below.
}
@Override
public Runnable herdrPollWait() {
return () -> {
throw new UnsupportedOperationException("healthy FakeHerdr must not poll");
};
}
}
private FleetdRuntime runtime;
private TestResourcePorts ports;
private Javalin boundApp;
@AfterEach
void tearDown() {
if (boundApp != null) {
boundApp.stop();
}
if (ports != null && ports.shutdownHook != null) {
ports.shutdownHook.run();
}
}
private static FleetConfig writeConfig(Path dir) throws Exception {
Path file = dir.resolve("fleetd.yaml");
Files.writeString(file, """
bind:
host: 127.0.0.1
port: 8765
auth:
mode: token
tokenEnv: %s
idleSleepGuard:
enabled: false
health:
enabled: true
notifications:
mode: webhook
profiles:
%s:
baseUrl: http://capacity.test:8000
model: test-model
maxLoad: 7
coordinator:
uri: amqp://fake-coordinator/vh
selfId: test-lead
peers:
- %s
""".formatted(TOKEN_ENV, PROFILE, PEER));
return FleetConfig.load(file);
}
private int assembleAndBind(Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
ports = new TestResourcePorts();
runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config,
new SubscriptionGuard(cfg.guard().hostSet())), ports);
boundApp = runtime.app().start("127.0.0.1", 0);
return boundApp.port();
}
private static String fleetList(int port) {
HttpRequest.Builder requestTemplate = HttpRequest.newBuilder()
.header("Authorization", "Bearer " + TOKEN);
McpClientTransport transport = HttpClientStreamableHttpTransport.builder("http://127.0.0.1:" + port)
.endpoint("/mcp")
.requestBuilder(requestTemplate)
.build();
try (McpSyncClient client = McpClient.sync(transport).build()) {
client.initialize();
McpSchema.CallToolResult response = client.callTool(McpSchema.CallToolRequest.builder("fleet_list")
.arguments(Map.of()).build());
return ((McpSchema.TextContent) response.content().getFirst()).text();
}
}
@Test
@Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD)
void fleetListReportsTheAssembledCapacitySource(@TempDir Path dir) throws Exception {
String response = fleetList(assembleAndBind(dir));
assertTrue(response.contains("\"capacity\""), response);
assertTrue(response.contains("\"" + PROFILE + "\""), response);
assertTrue(response.contains("\"maxLoad\":7"), response);
}
@Test
@Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD)
void fleetListReportsTheAssembledHealthCoverageSource(@TempDir Path dir) throws Exception {
String response = fleetList(assembleAndBind(dir));
assertTrue(response.contains("\"healthCoverage\":\"full\""), response);
}
@Test
@Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD)
void fleetListReportsTheAssembledCoordinatorPeers(@TempDir Path dir) throws Exception {
String response = fleetList(assembleAndBind(dir));
assertTrue(response.contains("\"coordinator\""), response);
assertTrue(response.contains("\"coordId\":\"" + PEER + "\""), response);
}
}
@@ -0,0 +1,368 @@
package dev.ltms.fleet;
import dev.ltms.fleet.config.ConfigRef;
import dev.ltms.fleet.config.FleetConfig;
import dev.ltms.fleet.guard.SubscriptionGuard;
import dev.ltms.fleet.herdr.FakeHerdr;
import dev.ltms.fleet.herdr.HerdrClient;
import dev.ltms.fleet.inject.CompletionResolver;
import dev.ltms.fleet.msg.Rendezvous;
import dev.ltms.fleet.msg.TurnToken;
import dev.ltms.fleet.session.MemberSession;
import io.javalin.Javalin;
import io.modelcontextprotocol.client.McpClient;
import io.modelcontextprotocol.client.McpSyncClient;
import io.modelcontextprotocol.client.transport.HttpClientStreamableHttpTransport;
import io.modelcontextprotocol.spec.McpClientTransport;
import io.modelcontextprotocol.spec.McpSchema;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.Timeout;
import org.junit.jupiter.api.io.TempDir;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.nio.file.Files;
import java.nio.file.Path;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.LongSupplier;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* fleetd #612 Shape A, unit r4 — {@link FleetdAssembly}'s {@code quarantineSource} (lines 471-472)
* and {@code outageSource} (lines 473-476 at {@code main} = {@code 141ae3b}), EACH of which feeds
* two separate consumers: {@code FleetMcp} ({@code fleet_profiles}) and {@code FleetApp}
* ({@code GET /profiles}), one call site ({@code :484}/{@code :486}) into the MCP constructor and
* the SAME shared local again ({@code :529}) into the REST constructor.
*
* <p><strong>The property pinned here</strong> (from the ticket): when the real assembly has built
* a real {@link dev.ltms.fleet.placement.BackendQuarantine} holding a quarantined credential, and a
* real {@link dev.ltms.fleet.placement.BackendOutagePolicy} holding a cooling-off credential, BOTH
* operator windows must report that state — the real assembled {@code FleetMcp} (reached through
* {@link FleetdRuntime#mcp()}) and the real assembled REST surface (reached through {@link
* FleetdRuntime#app()}). A mutation that starves one consumer while leaving the other wired must
* make only that consumer's assertion go red.
*
* <p><strong>No source-text assertion anywhere in this file.</strong> Both windows are read off the
* REAL running objects: {@code fleet_profiles} is called through a real MCP client over a real
* HTTP connection to the servlet {@link FleetdAssembly} actually mounted, and {@code GET /profiles}
* is called through a real {@link java.net.http.HttpClient} against the real bound {@link
* FleetdRuntime#app()}. Neither is a copy built alongside the assembly for this test's benefit.
*
* <p><strong>How this gets past CB-185's own pid-resolution dead end.</strong> {@code
* FleetdAssemblyFleetAppTest}'s class javadoc explains that {@code GET /sessions} cannot be driven
* over real HTTP here because {@code LsofPeerPidLookup} excludes its own pid and an in-process test
* client/server share one JVM pid — every such request resolves {@code ANONYMOUS} and is refused
* before the handler runs. {@code GET /profiles} and {@code fleet_profiles} sit behind the exact
* same {@code Authz.Action.READ} gate. This test sidesteps the dead end instead of hitting it:
* {@code auth.mode: token} (see {@link dev.ltms.fleet.auth.CallerResolver#resolve}) resolves a
* caller to {@code PRIMARY} from a valid {@code Authorization: Bearer} header ALONE, with no pid
* resolution involved at all — the same technique {@code FleetMcpContextExtractorTest} already uses
* to drive a real {@code fleet_whoami} call through the real transport.
*
* <p><strong>How the quarantined/cooling-off state is set up.</strong> Both {@code BackendQuarantine}
* and {@code BackendOutagePolicy} are private to the collaborators the assembly wires them into, and
* (unlike {@code quarantineSource()}) neither {@code FleetMcp} nor {@code FleetApp} exposes a public
* accessor for the live {@code BackendOutagePolicy} instance. Rather than add one (acceptance
* criterion 1: no production change), this test drives the REAL production classification path —
* exactly the recipe {@code FleetdExhaustedPatternAssemblyTest} (quarantine) and {@code
* FleetdCompletionResolverAssemblyTest} (cool-off) already proved works end to end against this same
* {@link FleetdAssembly#assembleAndStart}: acquire a real {@link MemberSession}, feed the real {@link
* CompletionResolver} a pane scrape matching the profile's configured {@code exhaustedPattern} /
* {@code errorPattern}, and let the real {@code exhaustionSink}/{@code backendErrorSink} write into
* the real, shared tracker. Each setup step asserts its own {@link Rendezvous.Kind} as a CONTROL —
* if the resolver were never actually exercised, the setup itself fails loudly before either window
* is ever read.
*/
class FleetdQuarantineOutageDualWindowAssemblyTest {
private static final String TOKEN = "s3cret-r4-token";
private static final String TOKEN_ENV = "FLEETD_R4_TEST_TOKEN";
private static final class ControllableResourcePorts implements ResourcePorts {
final FakeHerdr herdr;
final AtomicLong nowNanos = new AtomicLong(1_000_000_000L); // arbitrary non-zero start
Runnable shutdownHook;
ControllableResourcePorts(FakeHerdr herdr) {
this.herdr = herdr;
}
void advanceSeconds(long seconds) {
nowNanos.addAndGet(TimeUnit.SECONDS.toNanos(seconds));
}
@Override
public Map<String, String> environment() {
return Map.of(TOKEN_ENV, TOKEN);
}
@Override
public HerdrClient connectHerdr(Path socketPath) {
return herdr;
}
@Override
public Fleetd.AmqpOpener replyInboxOpener() {
// Never invoked: this test's config has no `broker:` block.
return (uri, prefetch) -> {
throw new UnsupportedOperationException("replyInboxOpener must not be called — no broker: block");
};
}
@Override
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
// Never invoked: this test's config has no `coordinator:` block.
return (uri, selfCoordId, prefetch) -> {
throw new UnsupportedOperationException(
"leadMailboxOpener must not be called — no coordinator: block is configured");
};
}
@Override
public LongSupplier nanoClock() {
return nowNanos::get;
}
@Override
public LongSupplier wallClockNanos() {
return nowNanos::get;
}
@Override
public ScheduledExecutorService newScheduler(String purpose) {
return Executors.newSingleThreadScheduledExecutor();
}
@Override
public void addShutdownHook(Runnable hook) {
this.shutdownHook = hook;
}
@Override
public void startHttp(Javalin app, String host, int port) {
// Deliberately never bind here — this test binds runtime.app() itself, for real, below.
}
@Override
public Runnable herdrPollWait() {
// Never invoked: this test's FakeHerdr answers immediately, so awaitHerdr never polls.
return () -> {
throw new UnsupportedOperationException("herdrPollWait must not be called — herdr is healthy");
};
}
}
private FleetdRuntime runtime;
private ControllableResourcePorts ports;
private Javalin boundApp;
@AfterEach
void tearDown() {
if (boundApp != null) {
boundApp.stop();
}
if (ports != null && ports.shutdownHook != null) {
ports.shutdownHook.run();
}
}
private static FleetConfig writeConfig(Path dir, int cooldownSeconds) throws Exception {
Path f = dir.resolve("fleetd.yaml");
Files.writeString(f, """
bind:
host: 127.0.0.1
port: 8765
auth:
mode: token
tokenEnv: %s
idleSleepGuard:
enabled: false
quarantineCooldownSeconds: %d
profiles:
exhaustprofile:
baseUrl: http://exhausthost.local:8000
model: sonnet
exhaustedPattern: "usage limit reached"
coolprofile:
baseUrl: http://coolhost.local:8000
model: sonnet
errorPattern: "credential outage"
guard:
offSubscriptionHosts:
- exhausthost.local
- coolhost.local
""".formatted(TOKEN_ENV, cooldownSeconds));
return FleetConfig.load(f);
}
/** Assembles the real graph, then binds the real {@code Javalin app} to an ephemeral port. */
private int assembleAndBind(Path dir) throws Exception {
FleetConfig cfg = writeConfig(dir, 120);
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
ports = new ControllableResourcePorts(new FakeHerdr());
runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
boundApp = runtime.app().start("127.0.0.1", 0);
return boundApp.port();
}
/**
* Drives the real assembled {@link CompletionResolver} through a scrape matching {@code
* exhaustprofile}'s configured {@code exhaustedPattern}, exactly {@code
* FleetdExhaustedPatternAssemblyTest}'s own recipe, so the real {@code exhaustionSink} quarantines
* the credential ({@code effectiveCredentialId() == "exhaustprofile"}, no explicit credentialId
* configured).
*/
private void quarantineExhaustProfile(Path dir) {
MemberSession session = runtime.sessions().acquire("exhaustprofile", null, dir.toString(), null);
String target = session.terminalId();
CompletionResolver completion = runtime.completion();
CompletableFuture<Rendezvous.Resolution> waiter = new CompletableFuture<>();
ports.herdr.readText("idle, nothing yet");
completion.onDelivered(target, new TurnToken(target, waiter, null));
ports.herdr.readText("usage limit reached: try again in a few hours");
ports.advanceSeconds(3); // clear CompletionResolver.MIN_TURN_NANOS (2s), no real sleep
completion.resolveBeforePostAction(target);
Rendezvous.Resolution resolution = waiter.getNow(null);
assertTrue(resolution != null && resolution.kind() == Rendezvous.Kind.BACKEND_EXHAUSTED,
"CONTROL: setup must classify as BACKEND_EXHAUSTED before either window is read — "
+ "if this fails, the assembled resolver was never actually exercised: " + resolution);
}
/**
* Drives the real assembled {@link CompletionResolver} with TWO distinct targets on {@code
* coolprofile}, each matching its configured {@code errorPattern}, exactly {@code
* FleetdCompletionResolverAssemblyTest}'s own recipe, so the real {@code backendErrorSink} cools
* the credential off ({@code effectiveCredentialId() == "coolprofile"}).
*/
private void coolOffCoolProfile(Path dir) {
MemberSession s1 = runtime.sessions().acquire("coolprofile", null, dir.toString(), null);
MemberSession s2 = runtime.sessions().acquire("coolprofile", null, dir.toString(), null);
CompletionResolver completion = runtime.completion();
String t1 = s1.terminalId();
CompletableFuture<Rendezvous.Resolution> w1 = new CompletableFuture<>();
ports.herdr.readText("idle 1");
completion.onDelivered(t1, new TurnToken(t1, w1, null));
ports.herdr.readText("credential outage: upstream 503");
ports.advanceSeconds(3);
completion.resolveBeforePostAction(t1);
Rendezvous.Resolution r1 = w1.getNow(null);
assertTrue(r1 != null && r1.kind() == Rendezvous.Kind.FAILED,
"CONTROL: target1's setup must classify FAILED (backend error): " + r1);
String t2 = s2.terminalId();
CompletableFuture<Rendezvous.Resolution> w2 = new CompletableFuture<>();
ports.herdr.readText("idle 2");
completion.onDelivered(t2, new TurnToken(t2, w2, null));
ports.herdr.readText("credential outage: upstream 503 again");
ports.advanceSeconds(3);
completion.resolveBeforePostAction(t2);
Rendezvous.Resolution r2 = w2.getNow(null);
assertTrue(r2 != null && r2.kind() == Rendezvous.Kind.FAILED,
"CONTROL: target2's setup must classify FAILED (backend error) — two distinct "
+ "targets are required to cross BackendOutagePolicy.THRESHOLD: " + r2);
}
/** Calls the real {@code fleet_profiles} tool over a real MCP client, token-authenticated as PRIMARY. */
private static McpSchema.CallToolResult callProfilesViaMcp(int port) {
HttpRequest.Builder requestTemplate = HttpRequest.newBuilder().header("Authorization", "Bearer " + TOKEN);
McpClientTransport transport = HttpClientStreamableHttpTransport.builder("http://127.0.0.1:" + port)
.endpoint("/mcp")
.requestBuilder(requestTemplate)
.build();
try (McpSyncClient client = McpClient.sync(transport).build()) {
client.initialize();
return client.callTool(McpSchema.CallToolRequest.builder("fleet_profiles").arguments(Map.of()).build());
}
}
private static String textOf(McpSchema.CallToolResult r) {
return ((McpSchema.TextContent) r.content().getFirst()).text();
}
/** Calls the real {@code GET /profiles} route over a real {@link HttpClient}, same token. */
private static String getProfilesViaRest(int port) throws Exception {
HttpClient http = HttpClient.newHttpClient();
HttpRequest req = HttpRequest.newBuilder(URI.create("http://127.0.0.1:" + port + "/profiles"))
.header("Authorization", "Bearer " + TOKEN)
.GET().build();
HttpResponse<String> res = http.send(req, HttpResponse.BodyHandlers.ofString());
assertEquals(200, res.statusCode(), "GET /profiles must succeed with the real token: " + res.body());
return res.body();
}
// --- quarantineSource (FleetdAssembly.java :471-472) --------------------------------------
@Test
@Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD)
void quarantinedCredentialIsReportedByTheRealAssembledFleetMcp(@TempDir Path dir) throws Exception {
int port = assembleAndBind(dir);
quarantineExhaustProfile(dir);
String out = textOf(callProfilesViaMcp(port));
assertTrue(out.contains("\"quarantined\""), "fleet_profiles must report a quarantined "
+ "section once the real BackendQuarantine holds a quarantined credential: " + out);
assertTrue(out.contains("\"exhaustprofile\""), out);
assertTrue(out.contains("\"quarantinedForSeconds\""), out);
}
@Test
@Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD)
void quarantinedCredentialIsReportedByTheRealAssembledFleetApp(@TempDir Path dir) throws Exception {
int port = assembleAndBind(dir);
quarantineExhaustProfile(dir);
String out = getProfilesViaRest(port);
assertTrue(out.contains("\"quarantined\""), "GET /profiles must report a quarantined "
+ "section once the real BackendQuarantine holds a quarantined credential: " + out);
assertTrue(out.contains("\"exhaustprofile\""), out);
assertTrue(out.contains("\"quarantinedForSeconds\""), out);
}
// --- outageSource (FleetdAssembly.java :473-476) -------------------------------------------
@Test
@Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD)
void coolingOffCredentialIsReportedByTheRealAssembledFleetMcp(@TempDir Path dir) throws Exception {
int port = assembleAndBind(dir);
coolOffCoolProfile(dir);
String out = textOf(callProfilesViaMcp(port));
assertTrue(out.contains("\"coolingOff\""), "fleet_profiles must report a coolingOff "
+ "section once the real BackendOutagePolicy holds a cooling-off credential: " + out);
assertTrue(out.contains("\"coolprofile\""), out);
assertTrue(out.contains("\"coolingOffForSeconds\""), out);
}
@Test
@Timeout(value = 15, unit = TimeUnit.SECONDS, threadMode = Timeout.ThreadMode.SEPARATE_THREAD)
void coolingOffCredentialIsReportedByTheRealAssembledFleetApp(@TempDir Path dir) throws Exception {
int port = assembleAndBind(dir);
coolOffCoolProfile(dir);
String out = getProfilesViaRest(port);
assertTrue(out.contains("\"coolingOff\""), "GET /profiles must report a coolingOff "
+ "section once the real BackendOutagePolicy holds a cooling-off credential: " + out);
assertTrue(out.contains("\"coolprofile\""), out);
assertTrue(out.contains("\"coolingOffForSeconds\""), out);
}
}