|
|
|
@@ -0,0 +1,307 @@
|
|
|
|
|
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.placement.PlacementException;
|
|
|
|
|
import dev.ltms.fleet.session.MemberSession;
|
|
|
|
|
import dev.ltms.fleet.session.WorktreeRequest;
|
|
|
|
|
import io.javalin.Javalin;
|
|
|
|
|
import org.junit.jupiter.api.Test;
|
|
|
|
|
import org.junit.jupiter.api.io.TempDir;
|
|
|
|
|
|
|
|
|
|
import java.nio.file.Files;
|
|
|
|
|
import java.nio.file.Path;
|
|
|
|
|
import java.util.List;
|
|
|
|
|
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.assertThrows;
|
|
|
|
|
import static org.junit.jupiter.api.Assertions.assertTrue;
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* fleetd #612 step 2, Unit B1 — replaces {@code FleetdCompletionResolverWiringTest} (deleted in
|
|
|
|
|
* this same commit), whose four tests read {@code Fleetd.java}'s source text and asserted the
|
|
|
|
|
* {@code CompletionResolver} construction call still named the right arguments. That proved the
|
|
|
|
|
* call site's spelling, never that the assembled resolver actually behaves differently when an
|
|
|
|
|
* argument is dropped.
|
|
|
|
|
*
|
|
|
|
|
* <p>These tests drive {@link FleetdAssembly#assembleAndStart} — the real boot composition,
|
|
|
|
|
* fleetd #612 Unit A — and read {@link FleetdRuntime#completion()}: the exact {@link
|
|
|
|
|
* CompletionResolver} instance the assembled daemon uses, never a copy built alongside it for the
|
|
|
|
|
* test's benefit. Two behaviours are pinned, matching the ticket's own two measured mutations:
|
|
|
|
|
*
|
|
|
|
|
* <ul>
|
|
|
|
|
* <li>the 8th constructor argument ({@code Fleetd.worktreeBranchLookup(sessions::roster)}) —
|
|
|
|
|
* {@link #assembledResolverReportsTheMembersWorktreeAndBranchInAFallbackReport}; and</li>
|
|
|
|
|
* <li>the 5th/6th arguments ({@code backendErrorPatterns}, {@code backendErrorSink}, both
|
|
|
|
|
* assigned from {@code Fleetd}'s extracted factories rather than an inline lambda) —
|
|
|
|
|
* {@link #assembledResolverClassifiesAndCoolsOffOnAConfiguredBackendErrorPattern}.</li>
|
|
|
|
|
* </ul>
|
|
|
|
|
*
|
|
|
|
|
* <p>Both tests bypass {@link dev.ltms.fleet.inject.StatusPoller} and drive {@link
|
|
|
|
|
* CompletionResolver#onDelivered} / {@link CompletionResolver#resolveBeforePostAction} directly —
|
|
|
|
|
* the same public, synchronous entry points {@code CompletionResolverTest} uses — with a
|
|
|
|
|
* hand-built {@link CompletableFuture} waiter, so no real poller loop or herdr status poll is
|
|
|
|
|
* needed. The pane scrape comes from {@link FakeHerdr#readText}; the elapsed-time floor
|
|
|
|
|
* ({@code CompletionResolver.MIN_TURN_NANOS}) is controlled via a fake, advanceable {@link
|
|
|
|
|
* ResourcePorts#nanoClock()} rather than a real sleep.
|
|
|
|
|
*/
|
|
|
|
|
class FleetdCompletionResolverAssemblyTest {
|
|
|
|
|
|
|
|
|
|
/** Same shape as {@code FleetdAssemblyLifecycleTest}'s fake, plus a nanoClock this test can advance. */
|
|
|
|
|
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();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
public HerdrClient connectHerdr(Path socketPath) {
|
|
|
|
|
return herdr;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
public Fleetd.AmqpOpener replyInboxOpener() {
|
|
|
|
|
// Never invoked: this test's config has no `broker:` block, so Fleetd.selectReplyInbox
|
|
|
|
|
// returns the in-memory inbox before calling the opener at all.
|
|
|
|
|
return (uri, prefetch) -> {
|
|
|
|
|
throw new UnsupportedOperationException("replyInboxOpener must not be called — no broker: block");
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@Override
|
|
|
|
|
public Fleetd.LeadMailboxOpener leadMailboxOpener() {
|
|
|
|
|
// Never invoked: no `coordinator:` block configured either.
|
|
|
|
|
return (uri, selfCoordId, prefetch) -> {
|
|
|
|
|
throw new UnsupportedOperationException("leadMailboxOpener must not be called — no coordinator: block");
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@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 a real port.
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private static FleetConfig writeConfig(Path dir, String profilesYaml, String extraGuardHost,
|
|
|
|
|
String worktreeRootYamlLine) throws Exception {
|
|
|
|
|
Path f = dir.resolve("fleetd.yaml");
|
|
|
|
|
Files.writeString(f, """
|
|
|
|
|
bind:
|
|
|
|
|
host: 127.0.0.1
|
|
|
|
|
port: 8765
|
|
|
|
|
idleSleepGuard:
|
|
|
|
|
enabled: false
|
|
|
|
|
%s
|
|
|
|
|
profiles:
|
|
|
|
|
%s
|
|
|
|
|
guard:
|
|
|
|
|
offSubscriptionHosts:
|
|
|
|
|
- %s
|
|
|
|
|
""".formatted(worktreeRootYamlLine == null ? "" : worktreeRootYamlLine, profilesYaml, extraGuardHost));
|
|
|
|
|
return FleetConfig.load(f);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private static void gitQuiet(Path cwd, String... args) throws Exception {
|
|
|
|
|
List<String> cmd = new java.util.ArrayList<>(List.of("git"));
|
|
|
|
|
cmd.addAll(List.of(args));
|
|
|
|
|
Process p = new ProcessBuilder(cmd).directory(cwd.toFile()).redirectErrorStream(true).start();
|
|
|
|
|
String out = new String(p.getInputStream().readAllBytes());
|
|
|
|
|
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git timed out: git " + String.join(" ", args));
|
|
|
|
|
assertEquals(0, p.exitValue(), "git " + String.join(" ", args) + " failed:\n" + out);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private static Path initRepo(Path dir) throws Exception {
|
|
|
|
|
Files.createDirectories(dir);
|
|
|
|
|
gitQuiet(dir, "init", "-q", "-b", "main");
|
|
|
|
|
gitQuiet(dir, "config", "user.email", "test@example.invalid");
|
|
|
|
|
gitQuiet(dir, "config", "user.name", "Test");
|
|
|
|
|
Files.writeString(dir.resolve("README.md"), "seed\n");
|
|
|
|
|
gitQuiet(dir, "add", "README.md");
|
|
|
|
|
gitQuiet(dir, "commit", "-q", "-m", "seed");
|
|
|
|
|
return dir;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* fleetd #248's first measured mutation: replacing {@code CompletionResolver}'s 8th constructor
|
|
|
|
|
* argument with the inert {@code _ -> null} compiles clean and leaves every existing test green
|
|
|
|
|
* — it silently drops fleetd #241's fallback-report location. This drives the real assembled
|
|
|
|
|
* resolver through a member echoing its own injected brief back (no {@code fleet_reply}), which
|
|
|
|
|
* resolves via {@code noReportMessage(target)}, and proves the real member's {@code branch} —
|
|
|
|
|
* only obtainable via {@code Fleetd.worktreeBranchLookup(sessions::roster)} reading the real,
|
|
|
|
|
* worktree-provisioned {@link MemberSession} — appears in the reported text.
|
|
|
|
|
*/
|
|
|
|
|
@Test
|
|
|
|
|
void assembledResolverReportsTheMembersWorktreeAndBranchInAFallbackReport(@TempDir Path dir) throws Exception {
|
|
|
|
|
Path repo = initRepo(dir.resolve("repo"));
|
|
|
|
|
FleetConfig cfg = writeConfig(dir, """
|
|
|
|
|
wtprofile:
|
|
|
|
|
baseUrl: http://wthost.local:8000
|
|
|
|
|
model: sonnet
|
|
|
|
|
""", "wthost.local", "worktreeRoot: " + dir.resolve("wts"));
|
|
|
|
|
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
|
|
|
|
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
|
|
|
|
ControllableResourcePorts ports = new ControllableResourcePorts(new FakeHerdr());
|
|
|
|
|
|
|
|
|
|
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
|
|
|
|
try {
|
|
|
|
|
MemberSession session = runtime.sessions().acquire("wtprofile", repo.toString(), repo.toString(),
|
|
|
|
|
null, new WorktreeRequest("fleetd-612-b1", null));
|
|
|
|
|
String target = session.terminalId();
|
|
|
|
|
String branch = session.branch();
|
|
|
|
|
assertTrue(branch != null && branch.startsWith("worker/"),
|
|
|
|
|
"sanity: a worktree-provisioned session must carry a real branch, got: " + branch);
|
|
|
|
|
|
|
|
|
|
CompletionResolver completion = runtime.completion();
|
|
|
|
|
CompletableFuture<Rendezvous.Resolution> waiter = new CompletableFuture<>();
|
|
|
|
|
String echoedBrief = "z".repeat(450); // >= CompletionResolver.ECHO_MIN_CHARS normalised chars
|
|
|
|
|
|
|
|
|
|
ports.herdr.readText("idle, nothing yet");
|
|
|
|
|
completion.onDelivered(target, new TurnToken(target, waiter, echoedBrief));
|
|
|
|
|
|
|
|
|
|
ports.herdr.readText(echoedBrief); // the pane just echoes the injected brief back — no real report
|
|
|
|
|
ports.advanceSeconds(3); // clear CompletionResolver.MIN_TURN_NANOS (2s) without a real sleep
|
|
|
|
|
completion.resolveBeforePostAction(target);
|
|
|
|
|
|
|
|
|
|
Rendezvous.Resolution resolution = waiter.getNow(null);
|
|
|
|
|
assertTrue(resolution != null, "the waiter must have resolved synchronously");
|
|
|
|
|
assertEquals(Rendezvous.Kind.COMPLETION, resolution.kind());
|
|
|
|
|
assertTrue(resolution.text().contains(CompletionResolver.NO_REPORT_PREFIX),
|
|
|
|
|
"sanity: must have gone down the noReportMessage sub-path: " + resolution.text());
|
|
|
|
|
assertTrue(resolution.text().contains("branch=" + branch),
|
|
|
|
|
"the assembled resolver must report the member's real branch (fleetd #241 via "
|
|
|
|
|
+ "fleetd #248's worktreeBranchLookup wiring); got: " + resolution.text());
|
|
|
|
|
} finally {
|
|
|
|
|
if (ports.shutdownHook != null) ports.shutdownHook.run();
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* fleetd #248's second measured mutation, and fleetd#201 Unit 5's own gap: replacing {@code
|
|
|
|
|
* backendErrorPatterns}/{@code backendErrorSink} with {@code BackendErrorPatternLookup.legacy()}
|
|
|
|
|
* / {@code BackendErrorSink.none()} compiles clean and leaves every existing behavioural test
|
|
|
|
|
* green.
|
|
|
|
|
*
|
|
|
|
|
* <p>Classification proof: this test's profile configures {@code errorPattern: "credential
|
|
|
|
|
* outage"} — text the built-in {@code (?i)\bAPI Error\s*:} fallback ({@code legacy()}'s only
|
|
|
|
|
* behaviour) never matches. So a real {@code Fleetd.backendErrorPatternLookup(...)} wiring
|
|
|
|
|
* classifies the send as {@code FAILED}; {@code legacy()} would fall through to the plain
|
|
|
|
|
* completion path instead ({@code Kind.COMPLETION}).
|
|
|
|
|
*
|
|
|
|
|
* <p>Cool-off proof: two distinct targets on the same profile/credential each classified as a
|
|
|
|
|
* backend error inside the 60s window must cool the credential off ({@link
|
|
|
|
|
* dev.ltms.fleet.placement.BackendOutagePolicy}, fleetd#201 Unit 5) — observable two ways: (1)
|
|
|
|
|
* the real {@code Fleetd.backendErrorSink(...)} marks each session {@code BACKEND_ERROR} (only
|
|
|
|
|
* the real sink calls {@code sessions.onBackendError}; {@code BackendErrorSink.none()} never
|
|
|
|
|
* does), and (2) a third explicit-profile spawn attempt is refused with a {@link
|
|
|
|
|
* PlacementException} naming the cool-off — only reachable because the real sink's {@code
|
|
|
|
|
* outagePolicy.record(...)} call actually ran.
|
|
|
|
|
*/
|
|
|
|
|
@Test
|
|
|
|
|
void assembledResolverClassifiesAndCoolsOffOnAConfiguredBackendErrorPattern(@TempDir Path dir) throws Exception {
|
|
|
|
|
FleetConfig cfg = writeConfig(dir, """
|
|
|
|
|
coolprofile:
|
|
|
|
|
baseUrl: http://coolhost.local:8000
|
|
|
|
|
model: sonnet
|
|
|
|
|
errorPattern: "credential outage"
|
|
|
|
|
""", "coolhost.local", null);
|
|
|
|
|
ConfigRef config = new ConfigRef(dir.resolve("fleetd.yaml"), cfg);
|
|
|
|
|
SubscriptionGuard guard = new SubscriptionGuard(cfg.guard().hostSet());
|
|
|
|
|
ControllableResourcePorts ports = new ControllableResourcePorts(new FakeHerdr());
|
|
|
|
|
|
|
|
|
|
FleetdRuntime runtime = FleetdAssembly.assembleAndStart(new AssemblyInputs(cfg, config, guard), ports);
|
|
|
|
|
try {
|
|
|
|
|
MemberSession session1 = runtime.sessions().acquire("coolprofile", null, dir.toString(), null);
|
|
|
|
|
MemberSession session2 = runtime.sessions().acquire("coolprofile", null, dir.toString(), null);
|
|
|
|
|
String target1 = session1.terminalId();
|
|
|
|
|
String target2 = session2.terminalId();
|
|
|
|
|
assertTrue(!target1.equals(target2), "sanity: the two spawns must be distinct targets");
|
|
|
|
|
|
|
|
|
|
CompletionResolver completion = runtime.completion();
|
|
|
|
|
|
|
|
|
|
CompletableFuture<Rendezvous.Resolution> waiter1 = new CompletableFuture<>();
|
|
|
|
|
ports.herdr.readText("idle 1");
|
|
|
|
|
completion.onDelivered(target1, new TurnToken(target1, waiter1, null));
|
|
|
|
|
ports.herdr.readText("credential outage: upstream 503");
|
|
|
|
|
ports.advanceSeconds(3);
|
|
|
|
|
completion.resolveBeforePostAction(target1);
|
|
|
|
|
Rendezvous.Resolution resolution1 = waiter1.getNow(null);
|
|
|
|
|
assertTrue(resolution1 != null, "target1's waiter must have resolved synchronously");
|
|
|
|
|
assertEquals(Rendezvous.Kind.FAILED, resolution1.kind(),
|
|
|
|
|
"a configured errorPattern the built-in fallback never matches must classify as "
|
|
|
|
|
+ "a backend error, not a plain completion; got: " + resolution1);
|
|
|
|
|
assertTrue(resolution1.text().contains("credential outage: upstream 503"), resolution1.text());
|
|
|
|
|
|
|
|
|
|
CompletableFuture<Rendezvous.Resolution> waiter2 = new CompletableFuture<>();
|
|
|
|
|
ports.herdr.readText("idle 2");
|
|
|
|
|
completion.onDelivered(target2, new TurnToken(target2, waiter2, null));
|
|
|
|
|
ports.herdr.readText("credential outage: upstream 503 again");
|
|
|
|
|
ports.advanceSeconds(3);
|
|
|
|
|
completion.resolveBeforePostAction(target2);
|
|
|
|
|
Rendezvous.Resolution resolution2 = waiter2.getNow(null);
|
|
|
|
|
assertTrue(resolution2 != null, "target2's waiter must have resolved synchronously");
|
|
|
|
|
assertEquals(Rendezvous.Kind.FAILED, resolution2.kind());
|
|
|
|
|
|
|
|
|
|
List<MemberSession> roster = runtime.sessions().roster();
|
|
|
|
|
assertTrue(roster.stream().anyMatch(s -> target1.equals(s.terminalId())
|
|
|
|
|
&& s.state() == MemberSession.State.BACKEND_ERROR),
|
|
|
|
|
"the real backendErrorSink must have transitioned target1 to BACKEND_ERROR: " + roster);
|
|
|
|
|
assertTrue(roster.stream().anyMatch(s -> target2.equals(s.terminalId())
|
|
|
|
|
&& s.state() == MemberSession.State.BACKEND_ERROR),
|
|
|
|
|
"the real backendErrorSink must have transitioned target2 to BACKEND_ERROR: " + roster);
|
|
|
|
|
|
|
|
|
|
PlacementException coolOff = assertThrows(PlacementException.class,
|
|
|
|
|
() -> runtime.sessions().acquire("coolprofile", null, dir.toString(), null),
|
|
|
|
|
"two distinct targets classified within the 60s window must cool the credential "
|
|
|
|
|
+ "off (BackendOutagePolicy), refusing a third explicit-profile spawn");
|
|
|
|
|
assertTrue(coolOff.getMessage().contains("cooling off"), coolOff.getMessage());
|
|
|
|
|
} finally {
|
|
|
|
|
if (ports.shutdownHook != null) ports.shutdownHook.run();
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|