Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 5289eb509f | |||
| 703a05db41 |
@@ -994,7 +994,11 @@ public final class FleetMcp {
|
||||
if (isBlank(target) || isBlank(msgId)) {
|
||||
return error("target and msgId are required");
|
||||
}
|
||||
messages.ackReply(target, msgId);
|
||||
if (!messages.ackReply(target, msgId)) {
|
||||
return error(msgId + " is not in " + target + "'s reply inbox (wrong id, wrong target, "
|
||||
+ "or already acked). Held lead-to-lead (peer) mail cannot be acked this way — "
|
||||
+ "read it with fleet_poll{coordId}.");
|
||||
}
|
||||
return text("acknowledged " + msgId);
|
||||
}
|
||||
|
||||
@@ -1766,7 +1770,10 @@ public final class FleetMcp {
|
||||
+ "has processed a reply and wants to confirm it, leaving other pending replies "
|
||||
+ "in the inbox for later drain.",
|
||||
objectSchema(Map.of(
|
||||
"target", stringProp("Worker session id whose inbox to ack from"),
|
||||
"target", stringProp("Worker session id whose inbox to ack from. Must name a "
|
||||
+ "reply actually queued for it — an id in no inbox, or a coord-id "
|
||||
+ "(peer held mail, read with fleet_poll{coordId} instead), errors "
|
||||
+ "rather than reporting a false success"),
|
||||
"msgId", stringProp("The message id to acknowledge")),
|
||||
List.of("target", "msgId")));
|
||||
}
|
||||
|
||||
@@ -17,7 +17,6 @@ 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;
|
||||
|
||||
@@ -595,23 +594,6 @@ 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();
|
||||
|
||||
@@ -394,17 +394,17 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void ack(String target, String msgId) {
|
||||
public boolean ack(String target, String msgId) {
|
||||
var perTarget = held.get(target);
|
||||
if (perTarget == null || perTarget == RELEASED) {
|
||||
return;
|
||||
return false;
|
||||
}
|
||||
Held h;
|
||||
synchronized (perTarget) {
|
||||
h = perTarget.remove(msgId);
|
||||
}
|
||||
if (h == null) {
|
||||
return; // never held (or already acked) — no-op
|
||||
return false; // never held (or already acked) — no-op
|
||||
}
|
||||
try {
|
||||
synchronized (channelLock) {
|
||||
@@ -418,6 +418,7 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
}
|
||||
throw new IllegalStateException("cannot ack reply " + msgId + " on " + queueName(target), e);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
private DeliverCallback deliverCallback(String target) {
|
||||
|
||||
@@ -59,16 +59,17 @@ public final class InMemoryReplyInbox implements ReplyInbox {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void ack(String target, String msgId) {
|
||||
public boolean ack(String target, String msgId) {
|
||||
if (!owned.contains(target)) {
|
||||
return;
|
||||
return false;
|
||||
}
|
||||
var perTarget = store.get(target);
|
||||
if (perTarget != null) {
|
||||
//noinspection SynchronizationOnLocalVariableOrMethodParameter
|
||||
synchronized (perTarget) {
|
||||
perTarget.remove(msgId);
|
||||
}
|
||||
if (perTarget == null) {
|
||||
return false;
|
||||
}
|
||||
//noinspection SynchronizationOnLocalVariableOrMethodParameter
|
||||
synchronized (perTarget) {
|
||||
return perTarget.remove(msgId) != null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -826,9 +826,13 @@ public final class MessageService {
|
||||
/**
|
||||
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
|
||||
* so that a subsequent drain or peek no longer returns it.
|
||||
*
|
||||
* @return {@code true} if an entry was actually removed, {@code false} if {@code msgId} was not
|
||||
* in {@code target}'s inbox (wrong id, wrong target, or already acked). The caller —
|
||||
* {@link dev.ltms.fleet.mcp.FleetMcp#ack} — must not report success on {@code false}.
|
||||
*/
|
||||
public void ackReply(String target, String msgId) {
|
||||
inbox.ack(target, msgId);
|
||||
public boolean ackReply(String target, String msgId) {
|
||||
return inbox.ack(target, msgId);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -46,6 +46,13 @@ public interface ReplyInbox {
|
||||
/** Non-destructive snapshot of pending replies for {@code target} (FIFO), empty list if none. */
|
||||
List<InboxMessage> peek(String target);
|
||||
|
||||
/** Remove the reply {@code msgId} for {@code target} once the primary has taken it. No-op if absent. */
|
||||
void ack(String target, String msgId);
|
||||
/**
|
||||
* Remove the reply {@code msgId} for {@code target} once the primary has taken it.
|
||||
*
|
||||
* @return {@code true} if an entry was actually removed, {@code false} if there was nothing to
|
||||
* remove (unknown {@code target}, unowned {@code target}, or a {@code msgId} not held for
|
||||
* it). A {@code false} is not an error — acking a {@code target} this daemon does not own is
|
||||
* part of the normal contract, not a failure.
|
||||
*/
|
||||
boolean ack(String target, String msgId);
|
||||
}
|
||||
|
||||
@@ -253,27 +253,15 @@ 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>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.
|
||||
* <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.
|
||||
*
|
||||
* @throws IllegalArgumentException if the decision names an unknown profile
|
||||
*/
|
||||
PeerHandle spawn(SpawnRequest req, PlacementDecision decision);
|
||||
default PeerHandle spawn(SpawnRequest req, PlacementDecision decision) {
|
||||
return spawn(req.withProfile(decision.profile()));
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve the effective working directory for a spawn {@code req} without actually spawning.
|
||||
|
||||
@@ -20,7 +20,6 @@ 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;
|
||||
@@ -124,11 +123,6 @@ 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();
|
||||
|
||||
@@ -352,7 +352,7 @@ class FleetMcpTest {
|
||||
fail("a lead fleet_reply must not publish to the worker inbox");
|
||||
}
|
||||
@Override public List<InboxMessage> peek(String target) { return List.of(); }
|
||||
@Override public void ack(String target, String msgId) { }
|
||||
@Override public boolean ack(String target, String msgId) { return false; }
|
||||
};
|
||||
MessageService leadMessages = new MessageService(agents, new Injector(agents), new Rendezvous(),
|
||||
inboxThatRejectsPublishes);
|
||||
@@ -1386,6 +1386,10 @@ class FleetMcpTest {
|
||||
|
||||
@Test
|
||||
void bridgeAckReturnsConfirmationForValidArgs() {
|
||||
// fleet_ack only reports success for a msgId actually queued in the target's inbox
|
||||
// (fleetd #437) — publish one via the inbox directly rather than asserting on a
|
||||
// fabricated id nothing ever queued.
|
||||
inbox.publish("term_a", "msg-1", "queued reply");
|
||||
McpSchema.CallToolResult res = FleetMcp.ack(messages, "term_a", "msg-1");
|
||||
assertNotEquals(Boolean.TRUE, res.isError());
|
||||
assertTrue(textOf(res).contains("msg-1"), "response should mention the msgId");
|
||||
@@ -1398,6 +1402,26 @@ class FleetMcpTest {
|
||||
assertTrue(FleetMcp.ack(messages, " ", "msg-1").isError());
|
||||
}
|
||||
|
||||
@Test
|
||||
void bridgeAckOfAnIdInNoInboxIsAnError() {
|
||||
// fleetd #437: fleet_ack used to say "acknowledged <msgId>" for a message it never
|
||||
// touched, because nothing in the chain reported hit vs. miss. "never-queued" is in no
|
||||
// inbox at all, so this must error rather than claim success.
|
||||
McpSchema.CallToolResult res = FleetMcp.ack(messages, "term_a", "never-queued");
|
||||
assertTrue(res.isError());
|
||||
assertTrue(textOf(res).contains("never-queued"), textOf(res));
|
||||
}
|
||||
|
||||
@Test
|
||||
void bridgeAckOfACoordIdTargetIsAnErrorNamingFleetPoll() {
|
||||
// A coord-id names a peer lead's held mailbox (LeadChannel/LeadMailbox), never a
|
||||
// worker's ReplyInbox — fleet_ack has no route to it and must say so, pointing at
|
||||
// fleet_poll{coordId} instead of reporting a false "acknowledged".
|
||||
McpSchema.CallToolResult res = FleetMcp.ack(messages, "coord-some-peer", "msg-1");
|
||||
assertTrue(res.isError());
|
||||
assertTrue(textOf(res).contains("fleet_poll{coordId}"), textOf(res));
|
||||
}
|
||||
|
||||
@Test
|
||||
void bridgeAckRemovesSpecificReply() {
|
||||
// Queue a reply and capture its msgId.
|
||||
@@ -1411,8 +1435,18 @@ class FleetMcpTest {
|
||||
var peeked = messages.drainReplies("term_a");
|
||||
assertEquals(1, peeked.size(), "one fresh reply in the inbox");
|
||||
|
||||
// ackReply works (no-op since published with a different UUID, but callable).
|
||||
assertDoesNotThrow(() -> messages.ackReply("term_a", msgId));
|
||||
// fleetd #437: msgId was already drained above (a fresh UUID each publish), so it is no
|
||||
// longer in the inbox — ackReply must now report that miss instead of pretending to ack.
|
||||
assertFalse(messages.ackReply("term_a", msgId));
|
||||
}
|
||||
|
||||
@Test
|
||||
void bridgeAckRemovingARealQueuedReplyReportsSuccessAndRemovesIt() {
|
||||
// The worker path must not change behaviour: acking a reply that IS still in the inbox
|
||||
// still succeeds and still removes it (fleetd #437).
|
||||
inbox.publish("term_a", "real-1", "still queued");
|
||||
assertTrue(messages.ackReply("term_a", "real-1"), "ack of a real queued reply must report true");
|
||||
assertTrue(inbox.peek("term_a").isEmpty(), "the acked reply must be gone from the inbox");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -20,7 +20,6 @@ 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;
|
||||
@@ -1124,68 +1123,6 @@ 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();
|
||||
|
||||
@@ -16,6 +16,7 @@ import java.util.List;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
@@ -221,6 +222,31 @@ class AmqpReplyInboxContractTest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void ackReportsHitVsMissAgainstARealBroker() throws Exception {
|
||||
// fleetd #437: fleet_ack said "acknowledged <msgId>" for a message it never touched,
|
||||
// because ReplyInbox.ack() (void) could not tell a hit from a miss. Pin the fixed
|
||||
// boolean contract against a real broker — the adapter fleetd actually runs live.
|
||||
String target = "worker-ack-contract-" + System.nanoTime();
|
||||
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
|
||||
inbox.own(target);
|
||||
|
||||
// Never held for this target at all: must report false, not throw.
|
||||
assertFalse(inbox.ack(target, "never-held"),
|
||||
"acking a msgId never held for an owned target must report false");
|
||||
|
||||
// A real message: first ack removes it and reports true...
|
||||
inbox.publish(target, "m1", "ack me");
|
||||
assertEquals(1, awaitPeek(inbox, target).size(), "the published reply should be held");
|
||||
assertTrue(inbox.ack(target, "m1"), "acking a held reply must report true");
|
||||
assertTrue(inbox.peek(target).isEmpty(), "an acked reply is dropped");
|
||||
|
||||
// ...and the second ack of the SAME msgId has nothing left to remove: false.
|
||||
assertFalse(inbox.ack(target, "m1"),
|
||||
"acking the same msgId twice must report false the second time");
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void confirmedPublishDeliversNormally() throws Exception {
|
||||
String target = "worker-confirm-" + System.nanoTime();
|
||||
|
||||
@@ -41,20 +41,20 @@ class InMemoryReplyInboxTest {
|
||||
@Test
|
||||
void ackRemovesTheMessage() {
|
||||
inbox.publish("term_a", "m1", "hello");
|
||||
inbox.ack("term_a", "m1");
|
||||
assertTrue(inbox.ack("term_a", "m1"), "fleetd #437: ack of a real entry must report true");
|
||||
assertTrue(inbox.peek("term_a").isEmpty(), "after ack, the message is gone");
|
||||
}
|
||||
|
||||
@Test
|
||||
void ackForUnknownMsgIdIsNoOp() {
|
||||
inbox.publish("term_a", "m1", "hello");
|
||||
inbox.ack("term_a", "no-such-id"); // no-op
|
||||
assertFalse(inbox.ack("term_a", "no-such-id"), "fleetd #437: a miss must report false"); // no-op
|
||||
assertEquals(1, inbox.peek("term_a").size(), "the published message is still there");
|
||||
}
|
||||
|
||||
@Test
|
||||
void ackForUnknownTargetIsNoOp() {
|
||||
inbox.ack("no-such-target", "m1"); // no-op, should not throw
|
||||
assertFalse(inbox.ack("no-such-target", "m1"), "fleetd #437: a miss must report false"); // no-op, should not throw
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -167,7 +167,7 @@ class InMemoryReplyInboxTest {
|
||||
@Test
|
||||
void peekAndAckAreNoOpsForUnownedTarget() {
|
||||
assertTrue(inbox.peek("term_not_owned").isEmpty());
|
||||
inbox.ack("term_not_owned", "m1"); // no-op, should not throw
|
||||
assertFalse(inbox.ack("term_not_owned", "m1"), "fleetd #437: a miss must report false"); // no-op, should not throw
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -23,7 +23,6 @@ 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;
|
||||
@@ -1037,11 +1036,6 @@ 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();
|
||||
@@ -1603,11 +1597,6 @@ 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");
|
||||
@@ -1672,11 +1661,6 @@ 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();
|
||||
@@ -1883,11 +1867,6 @@ 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