Compare commits
10 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| cfebc575ea | |||
| 822327eed5 | |||
| e4c703a51a | |||
| 3f036b2a62 | |||
| 82fae94c55 | |||
| b1f34c2e6b | |||
| 3f807d9f1b | |||
| e70263062c | |||
| 1a1e586b62 | |||
| c16d118f09 |
@@ -1397,11 +1397,11 @@ public final class FleetMcp {
|
||||
* {@code mailbox.pending} counts only broker-<em>ready</em> messages; a held message is already
|
||||
* an unacked delivery sitting with this consumer, so the normal, healthy state of a blocked lead
|
||||
* is {@code "pending": 0} next to a non-empty {@code held[]} — which invites the false reading
|
||||
* "these are only in memory, a restart will lose them". They are not: {@code LeadMailbox}
|
||||
* consumes with manual ack, so held mail is a durable broker delivery. {@code heldCount} is the
|
||||
* honest second number beside {@code pending} ({@code held.size()}, not left for the reader to
|
||||
* count the array), and {@code heldDurable} states the fact in words rather than leaving
|
||||
* {@code pending} as the only number next to {@code held[]}.
|
||||
* "these are only in memory, a restart will lose them". {@code heldCount} is the honest second
|
||||
* number beside {@code pending} ({@code held.size()}, not left for the reader to count the
|
||||
* array). {@code heldDurable} comes straight from {@link LeadChannel#heldDurable}, which the
|
||||
* channel implementation derives from what it actually did when it declared and consumed its own
|
||||
* queue (fleetd #440) — this method never asserts the fact itself.
|
||||
*/
|
||||
private static Map<String, Object> coordinatorView(CoordinationSource coordination) {
|
||||
LeadChannel channel = coordination.leadChannel();
|
||||
@@ -1415,7 +1415,7 @@ public final class FleetMcp {
|
||||
row.put("configured", true);
|
||||
row.put("mailbox", mailboxView(probe(channel, selfId)));
|
||||
row.put("heldCount", held.size());
|
||||
row.put("heldDurable", true);
|
||||
row.put("heldDurable", channel.heldDurable());
|
||||
row.put("held", held.stream().map(FleetMcp::heldView).toList());
|
||||
row.put("peers", coordination.peers().stream().map(p -> peerView(channel, p)).toList());
|
||||
return row;
|
||||
|
||||
@@ -17,6 +17,7 @@ import dev.ltms.fleet.peer.PeerHandle;
|
||||
import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import dev.ltms.fleet.peer.PeerUnreachableException;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import dev.ltms.fleet.placement.PlacementDecision;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -594,6 +595,23 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
req.sessionName(), spawned.agentSessionId(), spawned.receipt());
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritDoc}
|
||||
*
|
||||
* <p>fleetd #450: re-enters {@link #spawn(SpawnRequest)} with {@code decision}'s profile named
|
||||
* explicitly. This is the re-entering form the interface javadoc describes for a launcher with
|
||||
* no placement concept of its own — an instance of this class spawns a single adapter's own
|
||||
* profile set by explicit name only ({@link #place}/{@link #defaultProfileFor} are unoverridden
|
||||
* here and just wrap {@link #defaultProfile()}); it does no quarantine/cool-off/maxLoad/model-off
|
||||
* filtering of its own to re-apply. That filtering lives one layer up, in {@code
|
||||
* CompositePeerLauncher}, which is the launcher that routes across more than one profile and
|
||||
* therefore overrides this method with the routing form instead.
|
||||
*/
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
return spawn(req.withProfile(decision.profile()));
|
||||
}
|
||||
|
||||
/** The herdr daemon that owns this launcher's pane coordinates. */
|
||||
public HerdrClient herdr() {
|
||||
return agents.herdr();
|
||||
|
||||
@@ -42,6 +42,18 @@ public interface LeadChannel {
|
||||
/** This daemon's own lead coordination id — the mailbox it owns, and the {@code from} it sends as. */
|
||||
String selfCoordId();
|
||||
|
||||
/**
|
||||
* Whether a message sitting in {@link #peek}'s held set (fetched but not yet {@link #ack}ed) is
|
||||
* still safe if this daemon crashes or restarts right now — the conclusion of two independent
|
||||
* facts about how this channel owns its own queue: the queue was declared <em>durable</em>, and
|
||||
* the consumer that filled {@code held} uses <em>manual ack</em>, so an unacked delivery is still
|
||||
* owned by the broker rather than only in this process's memory. Both must hold for {@code true};
|
||||
* an implementation must derive this from what it actually did when it declared and consumed its
|
||||
* queue, never return a literal — fleetd #440 found {@code FleetMcp}'s {@code heldDurable} field
|
||||
* doing exactly that, unable to ever report {@code false} even after the fact stopped being true.
|
||||
*/
|
||||
boolean heldDurable();
|
||||
|
||||
/**
|
||||
* A non-destructive look at {@code coordId}'s mailbox — does it exist, how many messages are
|
||||
* waiting on it, and how many consumers are attached — without owning, consuming, or otherwise
|
||||
|
||||
@@ -86,6 +86,12 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
|
||||
private final Object channelLock = new Object();
|
||||
/** msgId → held delivery, for this mailbox's own queue only (there is exactly one). */
|
||||
private final LinkedHashMap<String, Held> held = new LinkedHashMap<>();
|
||||
/**
|
||||
* fleetd #440: the answer to {@link #heldDurable()}, set once by {@link #own()} from the exact
|
||||
* booleans it passed to {@code queueDeclare}/{@code basicConsume} — never a separate literal that
|
||||
* could drift from what those calls actually did.
|
||||
*/
|
||||
private boolean heldDurable;
|
||||
/** Successful broker acks on this connection, retained only to make a repeated caller ack quiet. */
|
||||
private final LinkedHashMap<String, Boolean> recentlyAcked = new LinkedHashMap<>();
|
||||
/** Bounds {@link #recentlyAcked}: it is only an idempotency aid, never delivery state. */
|
||||
@@ -191,13 +197,22 @@ public final class LeadMailbox implements LeadChannel, AutoCloseable {
|
||||
/** Declare + consume this daemon's own {@code lead.<selfCoordId>.inbox}. Called once, at construction. */
|
||||
private void own() throws IOException {
|
||||
String queue = queueName(selfCoordId);
|
||||
boolean durableQueue = true; // durable, non-exclusive, keep on idle
|
||||
boolean autoAck = false; // manual ack
|
||||
synchronized (channelLock) {
|
||||
channel.queueDeclare(queue, true, false, false, null); // durable, non-exclusive, keep on idle
|
||||
channel.basicConsume(queue, false, deliverCallback(), _ -> { }); // autoAck=false: manual ack
|
||||
channel.queueDeclare(queue, durableQueue, false, false, null);
|
||||
channel.basicConsume(queue, autoAck, deliverCallback(), _ -> { });
|
||||
}
|
||||
// fleetd #440: held mail is durable only while both hold — a durable queue AND manual ack.
|
||||
this.heldDurable = durableQueue && !autoAck;
|
||||
log.debug("lead mailbox owns queue {} for coord-id {}", queue, selfCoordId);
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean heldDurable() {
|
||||
return heldDurable;
|
||||
}
|
||||
|
||||
/**
|
||||
* Publish {@code msg} to {@code toCoordId}'s mailbox and block until the broker's publisher
|
||||
* confirm for it lands. Does <em>not</em> imply owning or consuming {@code toCoordId}'s queue.
|
||||
|
||||
@@ -253,15 +253,27 @@ public interface PeerLauncher {
|
||||
* matching profile; passing a request that names a <em>different</em>, explicit profile than
|
||||
* the decision it is paired with is a caller bug this method does not attempt to detect.
|
||||
*
|
||||
* <p>Default implementation for a launcher with no placement concept of its own: delegates to
|
||||
* {@link #spawn(SpawnRequest)} with the decision's profile named explicitly — its only spawn
|
||||
* contract, since there is no separate routing path to honor.
|
||||
* <p>No default implementation (fleetd #450): the two correct bodies disagree on purpose, so an
|
||||
* implementer must choose one rather than silently inherit whichever this interface happened to
|
||||
* provide. An implementer with no placement concept of its own — spawns a single profile, e.g.
|
||||
* {@code HerdrPeerLauncher} — should delegate to {@link #spawn(SpawnRequest)} with the decision's
|
||||
* profile named explicitly, since there is no separate routing path to honor there: the
|
||||
* explicit-profile branch it re-enters and the routing branch {@link #place} would have used are
|
||||
* the same thing. <strong>A launcher that routes across more than one profile — the way {@code
|
||||
* CompositePeerLauncher} routes across every configured adapter — MUST NOT re-enter {@link
|
||||
* #spawn(SpawnRequest)}.</strong> Doing so re-applies that single-argument method's
|
||||
* explicit-profile checks ({@code enforceNotQuarantined}, {@code enforceNotCoolingOff}, {@code
|
||||
* enforceMaxLoad}, {@code enforceModelEnabled} in {@code CompositePeerLauncher}), which can
|
||||
* refuse the very profile {@link #place} just chose, if the underlying placement state moved in
|
||||
* the window between the {@link #place} call and this one — the exact window this method and
|
||||
* {@link PlacementDecision} exist to close (fleetd #444). Before #450 this was a {@code default}
|
||||
* method that only {@code CompositePeerLauncher} overrode; a future placement-doing launcher
|
||||
* could have inherited the re-entering body silently and never known. Making it abstract turns
|
||||
* that silent inheritance into a compile error.
|
||||
*
|
||||
* @throws IllegalArgumentException if the decision names an unknown profile
|
||||
*/
|
||||
default PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
return spawn(req.withProfile(decision.profile()));
|
||||
}
|
||||
PeerHandle spawn(SpawnRequest req, PlacementDecision decision);
|
||||
|
||||
/**
|
||||
* Resolve the effective working directory for a spawn {@code req} without actually spawning.
|
||||
|
||||
@@ -20,6 +20,7 @@ import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import dev.ltms.fleet.placement.BackendOutagePolicy;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementDecision;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import dev.ltms.fleet.session.MemberSession;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
@@ -123,6 +124,11 @@ class FleetdBackendErrorSinkTest {
|
||||
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
throw new UnsupportedOperationException("not reachable — this test never acquires a session");
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return Set.of();
|
||||
|
||||
@@ -0,0 +1,79 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import ch.qos.logback.classic.Logger;
|
||||
import ch.qos.logback.classic.Level;
|
||||
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||
import ch.qos.logback.core.read.ListAppender;
|
||||
import dev.ltms.fleet.config.FleetConfig;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* Proves {@link Fleetd#main(String[])} calls every startup report before validation aborts startup.
|
||||
* The invalid non-loopback bind makes {@link FleetConfig#validateAll()} throw before {@code main}
|
||||
* can open the herdr socket or bind a port. The fixture also triggers every report, so removing any
|
||||
* one call from {@code main} leaves its expected log line absent.
|
||||
*/
|
||||
class FleetdStartupReportTest {
|
||||
|
||||
private static Level originalLevel;
|
||||
|
||||
private static ListAppender<ILoggingEvent> attach() {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
|
||||
originalLevel = logger.getLevel();
|
||||
logger.setLevel(Level.INFO);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
return appender;
|
||||
}
|
||||
|
||||
private static void detach(ListAppender<ILoggingEvent> appender) {
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(Fleetd.class);
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(originalLevel);
|
||||
}
|
||||
|
||||
private static boolean contains(ListAppender<ILoggingEvent> appender, String fragment) {
|
||||
return appender.list.stream()
|
||||
.map(ILoggingEvent::getFormattedMessage)
|
||||
.anyMatch(message -> message.contains(fragment));
|
||||
}
|
||||
|
||||
@Test
|
||||
void mainReportsEveryStartupGapBeforeValidationAborts(@TempDir Path dir) throws Exception {
|
||||
Path config = dir.resolve("fleetd.yaml");
|
||||
Files.writeString(config, """
|
||||
bind:
|
||||
host: 0.0.0.0
|
||||
port: 8765
|
||||
profiles:
|
||||
worker:
|
||||
baseUrl: https://llm.ltms.dev/v1
|
||||
gitTokenEnv: GITEA_TOKEN
|
||||
""");
|
||||
|
||||
ListAppender<ILoggingEvent> appender = attach();
|
||||
try {
|
||||
assertThrows(IllegalStateException.class, () -> Fleetd.main(new String[]{config.toString()}));
|
||||
} finally {
|
||||
detach(appender);
|
||||
}
|
||||
|
||||
assertTrue(contains(appender, "startup git host GITEA_HOST:"),
|
||||
"Fleetd.main must report the git host shape");
|
||||
assertTrue(contains(appender, "member trust model: members run as the same OS user"),
|
||||
"Fleetd.main must report the member trust model");
|
||||
assertTrue(contains(appender, "memberCredentials: absent or empty"),
|
||||
"Fleetd.main must report an absent memberCredentials policy");
|
||||
assertTrue(contains(appender, "exhaustedPattern: profile(s) [worker] have no exhaustedPattern configured"),
|
||||
"Fleetd.main must report profiles without exhaustedPattern");
|
||||
}
|
||||
}
|
||||
@@ -734,6 +734,32 @@ class FleetMcpTest {
|
||||
"must state the durability fact, not leave pending as the only number next to held[]: " + out);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #440: {@code heldDurable} must be a derived fact, not a literal — so it can report
|
||||
* {@code false} when the channel behind it says held mail is not durable (a non-durable queue,
|
||||
* or a consumer running with {@code autoAck=true}). A test that only ever asserts {@code true}
|
||||
* repeats the defect this ticket fixes.
|
||||
*/
|
||||
@Test
|
||||
void listReportsHeldDurableFalseWhenTheChannelSaysMailIsNotDurable() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
SessionManager sessions = new SessionManager(workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
|
||||
FakeLeadChannel channel = new FakeLeadChannel("mac-opus")
|
||||
.withMailbox("mac-opus", LeadChannel.MailboxState.exists("mac-opus", 0, 1))
|
||||
.withHeldDurable(false)
|
||||
.hold(new LeadMessage("m1", "fleet01-lead", "mac-opus", "one"));
|
||||
|
||||
McpSchema.CallToolResult res = FleetMcp.listFleet(
|
||||
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), sessions, null,
|
||||
FleetMcp.CapacitySource.none(), new FleetMcp.HealthCoverageSource(() -> "off"),
|
||||
FleetMcp.QuarantineSource.none(), Map.of(), "",
|
||||
new FleetMcp.CoordinationSource(channel, List.of()));
|
||||
|
||||
String out = textOf(res);
|
||||
assertTrue(out.contains("\"heldDurable\":false"),
|
||||
"heldDurable must follow the channel, not a hardcoded true: " + out);
|
||||
}
|
||||
|
||||
// ── fleetd #421: a lead reads (never consumes) its own held peer mail ──────────────────────
|
||||
|
||||
@Test
|
||||
@@ -834,6 +860,9 @@ class FleetMcpTest {
|
||||
@Override
|
||||
public String selfCoordId() { return "mac-opus"; }
|
||||
|
||||
@Override
|
||||
public boolean heldDurable() { return true; }
|
||||
|
||||
@Override
|
||||
public MailboxState inspect(String coordId) {
|
||||
started.countDown();
|
||||
|
||||
@@ -20,6 +20,7 @@ import dev.ltms.fleet.peer.PeerUnreachableException;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import dev.ltms.fleet.placement.BackendOutagePolicy;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementDecision;
|
||||
import dev.ltms.fleet.placement.PlacementException;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import org.junit.jupiter.api.Test;
|
||||
@@ -1123,6 +1124,68 @@ class CompositePeerLauncherTest {
|
||||
assertEquals(0, adapter.spawnCount("b"), "routedProfileFor never spawns anything");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #444: {@link PlacementDecision} exists to close the window between {@link
|
||||
* CompositePeerLauncher#place} and {@link CompositePeerLauncher#spawn(SpawnRequest,
|
||||
* PlacementDecision)} — the placement state must be free to move in that window without the
|
||||
* held decision being re-checked against the new state. Every quarantine test above resolves
|
||||
* and spawns in one call, so none of them ever open that window; this test is the one that
|
||||
* does: "sol" is placed FIRST, while nothing is quarantined yet, and only THEN is its
|
||||
* credential quarantined, before the held decision is spawned.
|
||||
*
|
||||
* <p>This is the test that tells the real override apart from the alternative body the ticket
|
||||
* measured: routing {@code decision.profile()} straight to its adapter (the real override)
|
||||
* never re-runs {@code enforceNotQuarantined}, so the spawn against the held decision still
|
||||
* succeeds on sol. Re-entering {@code spawn(req.withProfile(decision.profile()))} instead
|
||||
* lands in the explicit-profile branch, which refuses a now-quarantined sol outright — before
|
||||
* this test existed, replacing the real override's body with that re-entering call left the
|
||||
* whole suite green.
|
||||
*/
|
||||
@Test
|
||||
void spawnHonorsAPlacementDecisionEvenAfterItsProfileIsQuarantinedInTheWindowAfterPlace() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Map<String, FleetConfig.Profile> profiles = ordered(
|
||||
"sol", stubWorker("sol", "shared-openai"),
|
||||
"b", stubWorker("b"));
|
||||
// The adapter's OWN fallback default is "b", deliberately different from the profile place()
|
||||
// decides ("sol") — see the note below on why this must not be "sol" too.
|
||||
StubLauncher adapter = new StubLauncher("claude", herdr, profiles, "b", Set.of());
|
||||
BackendQuarantine quarantine = new BackendQuarantine(() -> 0L, TimeUnit.MINUTES.toNanos(30));
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(List.of(adapter), "sol", profiles,
|
||||
PlacementPolicies.fixed(), _ -> 0, null, quarantine);
|
||||
|
||||
// 1. Resolve BEFORE anything is quarantined — sol (definition order first, fixed policy) wins.
|
||||
// composite's own defaultProfile ("sol", the constructor arg above) never enters this: the
|
||||
// pool poolFor(DEV) resolves to is never empty here, so place() only ever reads that field as
|
||||
// a fallback for an empty pool, which this test does not exercise.
|
||||
PlacementDecision decision = composite.place(MemberRole.DEV);
|
||||
assertEquals("sol", decision.profile(), "sanity: nothing is quarantined yet, so sol is placed");
|
||||
|
||||
// 2. Move the placement state IN THE WINDOW between place() and spawn() — sol's credential
|
||||
// is now quarantined. A fresh place()/spawn(req) pair would fall through to b instead; the
|
||||
// held decision must not be re-evaluated against this new state at all.
|
||||
quarantine.quarantine("shared-openai");
|
||||
|
||||
// 3. Spawn against the HELD decision, not a fresh resolve.
|
||||
SpawnRequest req = new SpawnRequest(null, null, null, null, null, MemberRole.DEV);
|
||||
PeerHandle handle = composite.spawn(req, decision);
|
||||
|
||||
assertEquals("sol", handle.profile(),
|
||||
"the decision from place() is honored even though sol is now quarantined");
|
||||
// A fixture whose adapter falls back to "sol" too would let an UNSTAMPED request (one
|
||||
// routed but never given req.withProfile("sol")) land on spawnCount("sol") == 1 by
|
||||
// COINCIDENCE, since StubLauncher.spawn falls back to its own defaultProfile whenever
|
||||
// req.profileName() is blank. Giving the adapter "b" as its fallback instead means only an
|
||||
// actually-stamped request can produce this count — an unstamped one would count against
|
||||
// "b" and this assertion would fail.
|
||||
assertEquals(1, adapter.spawnCount("sol"),
|
||||
"the request that reached the delegate actually carried sol as its profile "
|
||||
+ "(the adapter's own fallback default is 'b', so this can't happen by accident)");
|
||||
assertEquals(0, adapter.spawnCount("b"),
|
||||
"b must never be touched — neither as the decision's profile nor as an unstamped "
|
||||
+ "request's accidental fallback");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aQuarantineLiftsOnTheInjectedClockAndTheProfileBecomesSpawnableAgain() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
|
||||
@@ -28,11 +28,19 @@ public final class FakeLeadChannel implements LeadChannel {
|
||||
private volatile IllegalStateException publishFailure;
|
||||
/** Canned {@link #inspect} results by coord-id — absent for any coord-id not configured here. */
|
||||
private final Map<String, MailboxState> mailboxes = new ConcurrentHashMap<>();
|
||||
/** fleetd #440: matches {@link LeadMailbox}'s real default (durable queue + manual ack) unless overridden. */
|
||||
private volatile boolean heldDurable = true;
|
||||
|
||||
public FakeLeadChannel(String selfCoordId) {
|
||||
this.selfCoordId = selfCoordId;
|
||||
}
|
||||
|
||||
/** Make {@link #heldDurable()} report {@code durable} — the fleetd #440 seam for the false case. */
|
||||
public FakeLeadChannel withHeldDurable(boolean durable) {
|
||||
this.heldDurable = durable;
|
||||
return this;
|
||||
}
|
||||
|
||||
/** Make {@link #inspect(String)} return {@code state} for {@code coordId} instead of "absent". */
|
||||
public FakeLeadChannel withMailbox(String coordId, MailboxState state) {
|
||||
mailboxes.put(coordId, state);
|
||||
@@ -80,6 +88,11 @@ public final class FakeLeadChannel implements LeadChannel {
|
||||
return mailboxes.getOrDefault(coordId, MailboxState.absent(coordId));
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean heldDurable() {
|
||||
return heldDurable;
|
||||
}
|
||||
|
||||
public List<LeadMessage> published() {
|
||||
return List.copyOf(published);
|
||||
}
|
||||
|
||||
@@ -207,6 +207,21 @@ class LeadMailboxTest {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #440: {@code heldDurable()} must be derived from what {@link LeadMailbox#own} actually
|
||||
* did against the real broker — a durable queue declare plus a manual-ack consumer — not a
|
||||
* hardcoded literal. This is the mutation-sensitive test: flip {@code own()}'s {@code autoAck}
|
||||
* local to {@code true} (or its {@code durableQueue} local to {@code false}) and this must fail.
|
||||
*/
|
||||
@Test
|
||||
void heldDurableReportsTrueBecauseTheQueueIsDurableAndTheConsumeIsManualAck() throws Exception {
|
||||
String self = coordId("lead-held-durable");
|
||||
try (LeadMailbox mailbox = LeadMailbox.open(uri(), self)) {
|
||||
assertTrue(mailbox.heldDurable(),
|
||||
"own() declares a durable queue and consumes with autoAck=false, so held mail is durable");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void inspectReportsAMissingMailboxAsAbsentRatherThanThrowing() throws Exception {
|
||||
String nobody = coordId("lead-inspect-nobody");
|
||||
|
||||
@@ -23,6 +23,7 @@ import dev.ltms.fleet.peer.PeerLauncher;
|
||||
import dev.ltms.fleet.peer.PeerUnreachableException;
|
||||
import dev.ltms.fleet.peer.SpawnRequest;
|
||||
import dev.ltms.fleet.placement.BackendQuarantine;
|
||||
import dev.ltms.fleet.placement.PlacementDecision;
|
||||
import dev.ltms.fleet.placement.PlacementPolicies;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
@@ -1036,6 +1037,11 @@ class SessionManagerTest {
|
||||
return delegate.spawn(req);
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
return delegate.spawn(req, decision);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return delegate.profiles();
|
||||
@@ -1597,6 +1603,11 @@ class SessionManagerTest {
|
||||
throw new UnsupportedOperationException("not reachable — the capability check refuses first");
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
throw new UnsupportedOperationException("not reachable — the capability check refuses first");
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return Set.of("stub-profile");
|
||||
@@ -1661,6 +1672,11 @@ class SessionManagerTest {
|
||||
return delegate.spawn(req);
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
return delegate.spawn(req, decision);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return delegate.profiles();
|
||||
@@ -1867,6 +1883,11 @@ class SessionManagerTest {
|
||||
return handle;
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
return spawn(req.withProfile(decision.profile()));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return Set.of("lazy");
|
||||
|
||||
Reference in New Issue
Block a user