Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 886ce1521d | |||
| d41aff4012 | |||
| 70a735b638 | |||
| 803c91ea6c | |||
| d438a74575 | |||
| 7467ffa252 | |||
| 6ab3a81af7 |
@@ -146,8 +146,9 @@ public record Principal(Role role, String terminal, long pid, String name) {
|
||||
}
|
||||
|
||||
/**
|
||||
* Stable identity used to own tickets and open turns. The unnamed primary has no owner key so
|
||||
* it can use the message layer's primary-wide ticket access rule.
|
||||
* Stable identity used to own tickets and open turns. The unnamed primary has no owner key —
|
||||
* {@code null} — and that is matched against a ticket's recorded owner the same way any other
|
||||
* key is: it owns a ticket another unnamed primary created, and nothing else.
|
||||
*/
|
||||
public String ownerKey() {
|
||||
return switch (role) {
|
||||
|
||||
@@ -93,7 +93,7 @@ import java.util.function.Supplier;
|
||||
*
|
||||
* <p><strong>Identity is resolved by the caller, never looked up here — a second fleetd #480
|
||||
* correction.</strong> The first version resolved the pane to clear via {@code
|
||||
* PrimaryRegistry#primaryTerminal()}. That is correct for a background loop with no caller (see
|
||||
* PrimaryRegistry#currentPrimaryTerminal()}. That is correct for a background loop with no caller (see
|
||||
* {@code dev.ltms.fleet.msg.LeadHeartbeatLoop}), but wrong here and a violation of this project's
|
||||
* own charter invariant 3 — "identity comes from the connection, never an argument." This daemon
|
||||
* can hold more than one labelled lead tab (see {@code LeadLauncher}'s fleetd #359 two-reading
|
||||
|
||||
@@ -600,7 +600,8 @@ public final class FleetMcp {
|
||||
(exchange, req) -> {
|
||||
McpSchema.CallToolResult denied = deny(exchange, toolAction("fleet_handover", req.arguments()), null);
|
||||
if (denied != null) return denied;
|
||||
return handover(leadRollover, callerTerminal(exchange), req.arguments());
|
||||
Principal caller = principal(exchange);
|
||||
return handover(leadRollover, messages, caller.terminal(), caller.ownerKey(), req.arguments());
|
||||
};
|
||||
|
||||
McpSchema.Tool fleetSend = sendTool();
|
||||
@@ -1358,9 +1359,10 @@ public final class FleetMcp {
|
||||
* {@code fleet_status}: the live lifecycle status of a worker session, plus — when the worker
|
||||
* is paused mid-turn in an async {@code fleet_ask} — the open question and how to answer it, so
|
||||
* a lead on its normal poll cadence does not need the ticket to notice. The question, its
|
||||
* {@code turnId} and its ticket id are shown only to the caller whose owner key created that
|
||||
* delegation, or to the unnamed primary; any other caller still sees the base status.
|
||||
* {@code callerOwner} comes from the calling connection's resolved principal.
|
||||
* {@code turnId} and its ticket id are shown only to the caller whose owner key matches the
|
||||
* delegation's creator — the unnamed primary matches only a delegation another unnamed
|
||||
* primary created; any other caller still sees the base status. {@code callerOwner} comes
|
||||
* from the calling connection's resolved principal.
|
||||
*/
|
||||
static McpSchema.CallToolResult status(MessageService messages, String sessionId, String callerOwner) {
|
||||
if (isBlank(sessionId)) {
|
||||
@@ -1481,14 +1483,15 @@ public final class FleetMcp {
|
||||
* #474 charter tool-surface gate), so every action here must degrade to a clean, structured
|
||||
* refusal naming {@code NOT_CONFIGURED} rather than ever throwing.
|
||||
*/
|
||||
static McpSchema.CallToolResult handover(LeadRollover leadRollover, String callerTerminal,
|
||||
static McpSchema.CallToolResult handover(LeadRollover leadRollover, MessageService messages,
|
||||
String callerTerminal, String callerOwner,
|
||||
Map<String, Object> args) {
|
||||
String action = str(args, "action");
|
||||
if (isBlank(action)) {
|
||||
return error("action is required: \"open\", \"confirm\", \"cancel\" or \"status\"");
|
||||
}
|
||||
return switch (action) {
|
||||
case "open" -> handoverOpen(leadRollover, callerTerminal, str(args, "reason"));
|
||||
case "open" -> handoverOpen(leadRollover, messages, callerTerminal, callerOwner, str(args, "reason"));
|
||||
case "confirm" -> handoverConfirm(leadRollover, callerTerminal, str(args, "token"),
|
||||
truthy(args, "operatorConfirmed"));
|
||||
case "cancel" -> handoverCancel(leadRollover, str(args, "token"));
|
||||
@@ -1504,7 +1507,8 @@ public final class FleetMcp {
|
||||
* IllegalStateException} for that), both degrade to the same clean {@code NOT_CONFIGURED}
|
||||
* refusal — never an escaping exception.
|
||||
*/
|
||||
private static McpSchema.CallToolResult handoverOpen(LeadRollover leadRollover, String callerTerminal,
|
||||
private static McpSchema.CallToolResult handoverOpen(LeadRollover leadRollover, MessageService messages,
|
||||
String callerTerminal, String callerOwner,
|
||||
String reason) {
|
||||
if (leadRollover == null) {
|
||||
return notConfigured();
|
||||
@@ -1522,6 +1526,9 @@ public final class FleetMcp {
|
||||
m.put("token", p.token());
|
||||
m.put("handoverPath", p.handoverPath());
|
||||
m.put("requestedAtMillis", p.requestedAtMillis());
|
||||
MessageService.Outstanding outstanding = messages.outstanding(callerOwner);
|
||||
m.put("outstandingTickets", outstanding.tickets());
|
||||
m.put("openAsks", outstanding.asks());
|
||||
return text(json(m));
|
||||
} catch (IllegalStateException e) {
|
||||
// leadRollover: was removed from config by a hot reload since this FleetMcp was
|
||||
@@ -2652,7 +2659,9 @@ public final class FleetMcp {
|
||||
+ "then use this to have fleetd end your pane's process and relaunch a fresh "
|
||||
+ "lead session bootstrapped against it. Four actions: 'open' (requests a "
|
||||
+ "token and the handoverPath you must write the handover file to before "
|
||||
+ "confirming), 'confirm' (validates every gate and — only if every one "
|
||||
+ "confirming — the response also lists outstandingTickets and openAsks, "
|
||||
+ "your own async delegations and fleet_ask turns, so their ids can go into "
|
||||
+ "the handover file too), 'confirm' (validates every gate and — only if every one "
|
||||
+ "passes — schedules the roll; it does NOT itself end your pane, the roll "
|
||||
+ "runs once this call's own turn ends), 'cancel' (drops a pending request "
|
||||
+ "without rolling), and 'status' (read-only: what happened to a token after "
|
||||
|
||||
@@ -12,6 +12,7 @@ import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Objects;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.CompletionException;
|
||||
@@ -1391,21 +1392,23 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #poll(String, String)}, with no caller owner key — the unnamed primary's ticket rule
|
||||
* checked, so this overload must only be used where the caller's identity is otherwise
|
||||
* irrelevant.
|
||||
* As {@link #poll(String, String)}, but bypasses the ownership check entirely via
|
||||
* {@link #INTERNAL_NO_OWNER_CHECK}. No production code calls this overload — it exists for
|
||||
* tests that only need the ticket's state and have no caller identity to pass.
|
||||
*/
|
||||
public TaskView poll(String ticket) {
|
||||
return poll(ticket, null);
|
||||
return poll(ticket, INTERNAL_NO_OWNER_CHECK);
|
||||
}
|
||||
|
||||
/**
|
||||
* Snapshot the state of an async delegation. Returns {@code null} for an unknown/expired ticket.
|
||||
* Refuses a {@code callerOwner} that differs from the owner that created the ticket (see
|
||||
* {@link #sendAsync(String, String, Runnable, Principal)}) with a {@link Phase#FAILED} view that
|
||||
* carries no reply text. The unnamed primary has a {@code null} owner key and is never refused.
|
||||
* Otherwise returns a {@link Phase#PENDING} view (with the live worker status as detail), a
|
||||
* {@link Phase#DONE} view carrying the reply, or a {@link Phase#FAILED} view with the reason.
|
||||
* carries no reply text. The unnamed primary's owner key is {@code null}, matched the same way
|
||||
* as any other key — it reads a ticket another unnamed primary created, and is refused on a
|
||||
* ticket a named caller created. Otherwise returns a {@link Phase#PENDING} view (with the live
|
||||
* worker status as detail), a {@link Phase#DONE} view carrying the reply, or a
|
||||
* {@link Phase#FAILED} view with the reason.
|
||||
*/
|
||||
public TaskView poll(String ticket, String callerOwner) {
|
||||
Task task = tasks.get(ticket);
|
||||
@@ -1452,13 +1455,27 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether {@code callerOwner} may read {@code task}'s state. A {@code null} caller key is the
|
||||
* unnamed primary and may read every ticket. Other callers must match the task's owner key. This
|
||||
* differs from {@link Rendezvous.Owner#permits}: a missing rendezvous owner is not an authenticated
|
||||
* Marker passed as {@code callerOwner} to bypass the ownership check entirely. No
|
||||
* {@link Principal#ownerKey()} ever produces this value — every real key is either
|
||||
* {@code null} (the unnamed primary) or prefixed with its role, such as {@code "worker:"} or
|
||||
* {@code "leader:"}. {@link #poll(String)} passes it; {@link #pendingAsk} has no matching
|
||||
* no-check overload, so this stays package-private for the test that drives the bypass
|
||||
* directly.
|
||||
*/
|
||||
static final String INTERNAL_NO_OWNER_CHECK = "internal:no-owner-check";
|
||||
|
||||
/**
|
||||
* Whether {@code callerOwner} may read {@code task}'s state. {@code callerOwner} is matched
|
||||
* against the task's recorded owner key by equality, including a {@code null} match — the
|
||||
* unnamed primary's owner key is {@code null}, so it owns a ticket another unnamed primary
|
||||
* created and nothing else, the same rule every other role follows. The only caller that
|
||||
* reads any ticket is {@link #INTERNAL_NO_OWNER_CHECK}. This differs from
|
||||
* {@link Rendezvous.Owner#permits}: a missing rendezvous owner is not an authenticated
|
||||
* unnamed primary, so that gate refuses every caller when no owner was recorded.
|
||||
*/
|
||||
private static boolean ownsTicket(Task task, String callerOwner) {
|
||||
return callerOwner == null || callerOwner.equals(task.creatorOwner);
|
||||
return INTERNAL_NO_OWNER_CHECK.equals(callerOwner)
|
||||
|| Objects.equals(callerOwner, task.creatorOwner);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1797,6 +1814,67 @@ public final class MessageService {
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* One ticket {@code callerOwner} created, still present in {@link #tasks}, surfaced by
|
||||
* {@link #outstanding} so a lead can carry its id into a handover file. {@link #phase} is the
|
||||
* same value {@link #poll} would report right now, terminal phases included: a {@code DONE} or
|
||||
* {@code FAILED} ticket stays in {@link #tasks} — and so stays reported here — until
|
||||
* {@link #pruneTerminalTickets} evicts it.
|
||||
*/
|
||||
public record OutstandingTicket(String ticket, Phase phase, String target) {
|
||||
}
|
||||
|
||||
/**
|
||||
* One worker session paused in {@code fleet_ask}, with the {@code turnId} that answers it,
|
||||
* surfaced by {@link #outstanding} alongside {@link OutstandingTicket}.
|
||||
*/
|
||||
public record OutstandingAsk(String ticket, String turnId, String workerSession) {
|
||||
}
|
||||
|
||||
/** The outstanding tickets and open asks a single call to {@link #outstanding} reports. */
|
||||
public record Outstanding(List<OutstandingTicket> tickets, List<OutstandingAsk> asks) {
|
||||
}
|
||||
|
||||
/**
|
||||
* Every ticket {@code callerOwner} created that is still in {@link #tasks} — including a
|
||||
* finished one nobody has polled yet, since {@link #pruneTerminalTickets} discards its reply
|
||||
* on a timer and a lead that does not carry its id forward can no longer read it after losing
|
||||
* its session's context — plus the subset of those whose worker is paused in
|
||||
* {@code fleet_ask}. Filtered by the same ownership rule as {@link #poll}:
|
||||
* {@link #ownsTicket(Task, String)}.
|
||||
*/
|
||||
public Outstanding outstanding(String callerOwner) {
|
||||
List<OutstandingTicket> tickets = new ArrayList<>();
|
||||
List<OutstandingAsk> asks = new ArrayList<>();
|
||||
for (Task task : tasks.values()) {
|
||||
if (!ownsTicket(task, callerOwner)) {
|
||||
continue;
|
||||
}
|
||||
Reply question = task.question;
|
||||
Phase phase;
|
||||
if (task.future.isDone()) {
|
||||
phase = terminalPhase(task.future);
|
||||
} else if (question != null) {
|
||||
phase = Phase.ASKING;
|
||||
asks.add(new OutstandingAsk(task.ticket, question.turnId(), task.target));
|
||||
} else {
|
||||
phase = Phase.PENDING;
|
||||
}
|
||||
tickets.add(new OutstandingTicket(task.ticket, phase, task.target));
|
||||
}
|
||||
return new Outstanding(tickets, asks);
|
||||
}
|
||||
|
||||
/** As {@link #poll}'s own terminal-result handling, reduced to just the {@link Phase}. */
|
||||
private static Phase terminalPhase(CompletableFuture<Reply> future) {
|
||||
try {
|
||||
Reply r = future.getNow(null);
|
||||
return r != null && r.completed() ? Phase.DONE : Phase.FAILED;
|
||||
} catch (CompletionException | java.util.concurrent.CancellationException e) {
|
||||
return Phase.FAILED;
|
||||
}
|
||||
}
|
||||
|
||||
/** Release the async executor. */
|
||||
public void close() {
|
||||
asyncExecutor.shutdown();
|
||||
|
||||
@@ -883,8 +883,7 @@ public final class FleetApp {
|
||||
body.put("ready", deliverable.test(id));
|
||||
// A worker paused mid-turn in an async fleet_ask is otherwise invisible to a status
|
||||
// poll — surface the open question and how to answer it, same as fleet_poll's
|
||||
// Phase.ASKING view, but only to the caller whose owner key created that delegation, or
|
||||
// to the unnamed primary.
|
||||
// Phase.ASKING view, but only to the caller whose owner key created that delegation.
|
||||
Principal caller = ctx.attribute(CALLER);
|
||||
MessageService.PendingAsk ask = messages.pendingAsk(id, caller == null ? null : caller.ownerKey());
|
||||
if (ask != null) {
|
||||
|
||||
@@ -56,9 +56,14 @@ class FleetMcpHandoverTest {
|
||||
|
||||
private static final String LEAD = "term_lead";
|
||||
private static final String OTHER_LEAD = "term_other_lead";
|
||||
private static final String LEAD_OWNER = "leader:lead";
|
||||
|
||||
private final FakeHerdr herdr = new FakeHerdr();
|
||||
private final AgentControl agents = new AgentControl(herdr);
|
||||
/** Fed to every direct {@code FleetMcp.handover} call below — none of this class's own tests
|
||||
* exercise ticket/ask ownership, so a single instance with no delegations is enough. */
|
||||
private final MessageService messages = new MessageService(agents, new Injector(agents),
|
||||
new Rendezvous(), new InMemoryReplyInbox());
|
||||
private FleetMcp mcp;
|
||||
|
||||
@AfterEach
|
||||
@@ -151,16 +156,16 @@ class FleetMcpHandoverTest {
|
||||
@DisplayName("with leadRollover: absent, every action returns a clean NOT_CONFIGURED refusal and never throws")
|
||||
void nullLeadRolloverRefusesCleanlyForEveryAction() {
|
||||
McpSchema.CallToolResult open = assertDoesNotThrow(
|
||||
() -> FleetMcp.handover(null, LEAD, Map.of("action", "open")));
|
||||
() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER, Map.of("action", "open")));
|
||||
assertFalse(open.isError(), "a refusal is not a protocol error: " + textOf(open));
|
||||
assertTrue(textOf(open).contains("NOT_CONFIGURED"), textOf(open));
|
||||
|
||||
McpSchema.CallToolResult confirm = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD,
|
||||
McpSchema.CallToolResult confirm = assertDoesNotThrow(() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER,
|
||||
Map.of("action", "confirm", "token", "whatever")));
|
||||
assertFalse(confirm.isError());
|
||||
assertTrue(textOf(confirm).contains("NOT_CONFIGURED"), textOf(confirm));
|
||||
|
||||
McpSchema.CallToolResult cancel = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD,
|
||||
McpSchema.CallToolResult cancel = assertDoesNotThrow(() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER,
|
||||
Map.of("action", "cancel", "token", "whatever")));
|
||||
assertFalse(cancel.isError());
|
||||
assertTrue(textOf(cancel).contains("NOT_CONFIGURED"), textOf(cancel));
|
||||
@@ -169,11 +174,11 @@ class FleetMcpHandoverTest {
|
||||
@Test
|
||||
@DisplayName("a blank/unknown action is a clean tool error, never an exception")
|
||||
void unknownActionIsACleanError() {
|
||||
McpSchema.CallToolResult missing = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD, Map.of()));
|
||||
McpSchema.CallToolResult missing = assertDoesNotThrow(() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER, Map.of()));
|
||||
assertTrue(missing.isError());
|
||||
|
||||
McpSchema.CallToolResult bogus = assertDoesNotThrow(
|
||||
() -> FleetMcp.handover(null, LEAD, Map.of("action", "bogus")));
|
||||
() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER, Map.of("action", "bogus")));
|
||||
assertTrue(bogus.isError());
|
||||
}
|
||||
|
||||
@@ -229,7 +234,7 @@ class FleetMcpHandoverTest {
|
||||
Files.writeString(handover, "not written yet");
|
||||
LeadRollover rollover = newRollover(handover.toString());
|
||||
|
||||
McpSchema.CallToolResult openResult = FleetMcp.handover(rollover, LEAD, Map.of("action", "open"));
|
||||
McpSchema.CallToolResult openResult = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER, Map.of("action", "open"));
|
||||
assertFalse(openResult.isError(), textOf(openResult));
|
||||
String token = extractToken(textOf(openResult));
|
||||
|
||||
@@ -238,13 +243,13 @@ class FleetMcpHandoverTest {
|
||||
Thread.sleep(50);
|
||||
Files.writeString(handover, "the real handover content");
|
||||
|
||||
McpSchema.CallToolResult wrongCaller = FleetMcp.handover(rollover, OTHER_LEAD,
|
||||
McpSchema.CallToolResult wrongCaller = FleetMcp.handover(rollover, messages, OTHER_LEAD, "leader:other-lead",
|
||||
Map.of("action", "confirm", "token", token));
|
||||
assertFalse(wrongCaller.isError(), "a refusal is a legitimate outcome, not a protocol error");
|
||||
assertTrue(textOf(wrongCaller).contains("NOT_YOUR_ROLLOVER"),
|
||||
"a different lead terminal confirming must surface NOT_YOUR_ROLLOVER: " + textOf(wrongCaller));
|
||||
|
||||
McpSchema.CallToolResult confirmed = FleetMcp.handover(rollover, LEAD,
|
||||
McpSchema.CallToolResult confirmed = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
|
||||
Map.of("action", "confirm", "token", token));
|
||||
assertFalse(confirmed.isError(), textOf(confirmed));
|
||||
assertTrue(textOf(confirmed).contains("\"accepted\":true"),
|
||||
@@ -258,7 +263,7 @@ class FleetMcpHandoverTest {
|
||||
void cancelUnknownTokenIsCleanNotAFailure() {
|
||||
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
|
||||
|
||||
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
|
||||
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
|
||||
Map.of("action", "cancel", "token", "does-not-exist"));
|
||||
assertFalse(r.isError());
|
||||
assertTrue(textOf(r).contains("\"cancelled\":false"), textOf(r));
|
||||
@@ -269,9 +274,9 @@ class FleetMcpHandoverTest {
|
||||
@DisplayName("cancel on a token actually opened reports cancelled:true")
|
||||
void cancelKnownTokenSucceeds() {
|
||||
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
|
||||
String token = extractToken(textOf(FleetMcp.handover(rollover, LEAD, Map.of("action", "open"))));
|
||||
String token = extractToken(textOf(FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER, Map.of("action", "open"))));
|
||||
|
||||
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
|
||||
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
|
||||
Map.of("action", "cancel", "token", token));
|
||||
assertFalse(r.isError());
|
||||
assertTrue(textOf(r).contains("\"cancelled\":true"), textOf(r));
|
||||
@@ -282,7 +287,7 @@ class FleetMcpHandoverTest {
|
||||
@Test
|
||||
@DisplayName("status on a null LeadRollover is a clean NOT_CONFIGURED refusal, never a throw")
|
||||
void statusWithNullLeadRolloverRefusesCleanly() {
|
||||
McpSchema.CallToolResult r = assertDoesNotThrow(() -> FleetMcp.handover(null, LEAD,
|
||||
McpSchema.CallToolResult r = assertDoesNotThrow(() -> FleetMcp.handover(null, messages, LEAD, LEAD_OWNER,
|
||||
Map.of("action", "status", "token", "whatever")));
|
||||
assertFalse(r.isError());
|
||||
assertTrue(textOf(r).contains("NOT_CONFIGURED"), textOf(r));
|
||||
@@ -293,7 +298,7 @@ class FleetMcpHandoverTest {
|
||||
void statusOnUnknownTokenReportsUnknown() {
|
||||
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
|
||||
|
||||
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
|
||||
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
|
||||
Map.of("action", "status", "token", "does-not-exist"));
|
||||
assertFalse(r.isError());
|
||||
assertTrue(textOf(r).contains("\"state\":\"UNKNOWN\""), textOf(r));
|
||||
@@ -303,9 +308,9 @@ class FleetMcpHandoverTest {
|
||||
@DisplayName("status on a token that is still pending (opened, not confirmed) reports PENDING")
|
||||
void statusOnPendingTokenReportsPending() {
|
||||
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
|
||||
String token = extractToken(textOf(FleetMcp.handover(rollover, LEAD, Map.of("action", "open"))));
|
||||
String token = extractToken(textOf(FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER, Map.of("action", "open"))));
|
||||
|
||||
McpSchema.CallToolResult r = FleetMcp.handover(rollover, LEAD,
|
||||
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
|
||||
Map.of("action", "status", "token", token));
|
||||
assertFalse(r.isError());
|
||||
assertTrue(textOf(r).contains("\"state\":\"PENDING\""), textOf(r));
|
||||
@@ -323,5 +328,40 @@ class FleetMcpHandoverTest {
|
||||
"the tool's own description must advertise the 'status' action: " + tool.description());
|
||||
}
|
||||
|
||||
// --- unit 5: open() reports outstanding tickets and open asks ------------------------------
|
||||
|
||||
@Test
|
||||
@DisplayName("open() keeps token/handoverPath/requestedAtMillis and reports empty outstanding collections when the caller has nothing")
|
||||
void openReportsEmptyOutstandingCollectionsWhenCallerHasNothing() {
|
||||
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
|
||||
|
||||
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
|
||||
Map.of("action", "open"));
|
||||
assertFalse(r.isError(), textOf(r));
|
||||
String json = textOf(r);
|
||||
assertTrue(json.contains("\"token\":"), json);
|
||||
assertTrue(json.contains("\"handoverPath\":"), json);
|
||||
assertTrue(json.contains("\"requestedAtMillis\":"), json);
|
||||
assertTrue(json.contains("\"outstandingTickets\":[]"),
|
||||
"a caller with nothing gets an empty array, not an absent key: " + json);
|
||||
assertTrue(json.contains("\"openAsks\":[]"),
|
||||
"a caller with nothing gets an empty array, not an absent key: " + json);
|
||||
}
|
||||
|
||||
@Test
|
||||
@DisplayName("open() reports an owned pending ticket with its phase and target")
|
||||
void openReportsAnOwnedPendingTicketWithItsPhase() {
|
||||
LeadRollover rollover = newRollover(tmp.resolve("h.md").toString());
|
||||
String ticket = messages.sendAsync("term_worker", "a task", null, Principal.leader("lead", LEAD, 1));
|
||||
|
||||
McpSchema.CallToolResult r = FleetMcp.handover(rollover, messages, LEAD, LEAD_OWNER,
|
||||
Map.of("action", "open"));
|
||||
assertFalse(r.isError(), textOf(r));
|
||||
String json = textOf(r);
|
||||
assertTrue(json.contains("\"ticket\":\"" + ticket + "\""), json);
|
||||
assertTrue(json.contains("\"phase\":\"PENDING\""), json);
|
||||
assertTrue(json.contains("\"target\":\"term_worker\""), json);
|
||||
}
|
||||
|
||||
// --- acceptance 7 (wiring) is covered by FleetdLeadRolloverWiringTest, unchanged -----------
|
||||
}
|
||||
|
||||
@@ -2013,8 +2013,10 @@ class FleetMcpTest {
|
||||
|
||||
/**
|
||||
* {@code fleet_status}'s pending-ask block (the question, its {@code turnId} and its ticket)
|
||||
* is shown only to the caller whose owner key created the delegation, or to the unnamed primary.
|
||||
* A different caller still sees the base status line, but none of the pending-ask fields.
|
||||
* is shown only to the caller whose owner key created the delegation. An unnamed primary is
|
||||
* held to the same rule: its owner key is {@code null}, which here does not match the named
|
||||
* worker that created the delegation, so it sees none of the pending-ask fields either — the
|
||||
* same as any other non-creating caller.
|
||||
*/
|
||||
@Test
|
||||
void statusGatesThePendingAskFieldsByTheDelegationsCreatorOwner() throws Exception {
|
||||
@@ -2054,8 +2056,13 @@ class FleetMcpTest {
|
||||
assertTrue(creatorStatus.contains(ticket), "the creator must see the ticket: " + creatorStatus);
|
||||
|
||||
String unnamed = textOf(FleetMcp.status(messages, T, null));
|
||||
assertTrue(unnamed.contains("which config file?"),
|
||||
"a caller with no terminal (the unnamed primary) must see the question: " + unnamed);
|
||||
assertTrue(unnamed.startsWith("idle"), "the base status must still be shown: " + unnamed);
|
||||
assertFalse(unnamed.contains("which config file?"),
|
||||
"an unnamed primary must not see a question on a delegation a named worker created: " + unnamed);
|
||||
assertFalse(unnamed.contains(asking.turnId()),
|
||||
"a non-creating unnamed primary must not see the turnId: " + unnamed);
|
||||
assertFalse(unnamed.contains(ticket),
|
||||
"a non-creating unnamed primary must not see the ticket: " + unnamed);
|
||||
|
||||
// Clean up the still-open ask so the background thread does not linger past the test.
|
||||
String turnId = asking.turnId();
|
||||
|
||||
@@ -948,7 +948,7 @@ class MessageServiceTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void unnamedPrimaryStillReadsAnyTicket() throws Exception {
|
||||
void unnamedPrimaryIsRefusedFromANamedLeadsTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task", null,
|
||||
Principal.leader("opus", "term_lead", 1));
|
||||
awaitWaiting();
|
||||
@@ -956,8 +956,51 @@ class MessageServiceTest {
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker works
|
||||
assertTrue(rendezvous.resolve(T, "primary-visible result"), "a reply resolves the async send");
|
||||
|
||||
MessageService.TaskView view = driveAsyncTicketToDone(ticket, null);
|
||||
assertNotNull(view, "the unnamed primary must be able to read any ticket");
|
||||
MessageService.TaskView refused = messages.poll(ticket, Principal.primary(1).ownerKey());
|
||||
assertNotNull(refused, "a different owner gets a refusal, not silence");
|
||||
assertEquals(MessageService.Phase.FAILED, refused.phase());
|
||||
assertEquals("forbidden: this ticket was created by a different session", refused.detail());
|
||||
assertNull(refused.reply(), "a refusal must never carry the reply text");
|
||||
assertFalse(String.valueOf(refused).contains("primary-visible result"),
|
||||
"the reply text must not appear anywhere in the refused view");
|
||||
}
|
||||
|
||||
/**
|
||||
* Positive control for {@link #unnamedPrimaryIsRefusedFromANamedLeadsTicket}: without this,
|
||||
* that test would pass just as well if {@code poll} refused every caller.
|
||||
*/
|
||||
@Test
|
||||
void unnamedPrimaryReadsItsOwnTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task", null, Principal.primary(1));
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker works
|
||||
assertTrue(rendezvous.resolve(T, "primary-visible result"), "a reply resolves the async send");
|
||||
|
||||
MessageService.TaskView view = driveAsyncTicketToDone(ticket, Principal.primary(2).ownerKey());
|
||||
assertNotNull(view, "an unnamed primary must be able to read a ticket another unnamed primary created");
|
||||
assertEquals(MessageService.Phase.DONE, view.phase());
|
||||
assertEquals("primary-visible result", view.reply());
|
||||
}
|
||||
|
||||
@Test
|
||||
void theOneArgPollOverloadBypassesOwnershipEntirely() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task", null,
|
||||
Principal.leader("opus", "term_lead", 1));
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker works
|
||||
assertTrue(rendezvous.resolve(T, "primary-visible result"), "a reply resolves the async send");
|
||||
|
||||
MessageService.TaskView view = null;
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (view == null || view.phase() != MessageService.Phase.DONE) {
|
||||
if (System.currentTimeMillis() >= deadline) break;
|
||||
view = messages.poll(ticket); // the one-arg, no-check overload -- no caller owner key at all
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertNotNull(view, "the internal bypass must read a ticket owned by a named lead");
|
||||
assertEquals(MessageService.Phase.DONE, view.phase());
|
||||
assertEquals("primary-visible result", view.reply());
|
||||
}
|
||||
@@ -1978,7 +2021,9 @@ class MessageServiceTest {
|
||||
|
||||
/**
|
||||
* A caller's owner key must match the key that created the delegation to see its pending
|
||||
* question. The unnamed primary always sees it.
|
||||
* question. The unnamed primary is held to the same rule as everyone else: its key is
|
||||
* {@code null}, which here does not match the named worker that created this delegation, so
|
||||
* it is refused too.
|
||||
*/
|
||||
@Test
|
||||
void pendingAskGatesTheQuestionByTheDelegationsCreatorOwner() throws Exception {
|
||||
@@ -1998,10 +2043,65 @@ class MessageServiceTest {
|
||||
assertNotNull(own, "the creating caller must see its own open question");
|
||||
assertEquals("which config file?", own.question());
|
||||
|
||||
MessageService.PendingAsk unnamed = messages.pendingAsk(T, null);
|
||||
assertNotNull(unnamed, "a caller with no terminal (the unnamed primary) must always see the question");
|
||||
assertNull(messages.pendingAsk(T, Principal.primary(1).ownerKey()),
|
||||
"an unnamed primary must not see a question on a delegation a named worker created");
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000, creator.ownerKey()));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
/**
|
||||
* Positive control for {@link #pendingAskGatesTheQuestionByTheDelegationsCreatorOwner}:
|
||||
* without this, that test's refusal would pass just as well if {@code pendingAsk} refused
|
||||
* every caller. Here the delegation's creator is itself an unnamed primary (owner key
|
||||
* {@code null}), so another unnamed primary's {@code null} key must still match it.
|
||||
*/
|
||||
@Test
|
||||
void unnamedPrimarySeesItsOwnDelegationsPendingQuestion() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks", null, Principal.primary(1));
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
||||
awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
MessageService.PendingAsk unnamed = messages.pendingAsk(T, Principal.primary(2).ownerKey());
|
||||
assertNotNull(unnamed, "an unnamed primary must see the question on a delegation another unnamed primary created");
|
||||
assertEquals("which config file?", unnamed.question());
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(unnamed.turnId(), "config.yaml", 5000, null));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code pendingAsk} has no public no-check overload the way {@link MessageService#poll}
|
||||
* does, so this drives {@link MessageService#INTERNAL_NO_OWNER_CHECK} directly — the only way
|
||||
* to exercise the bypass for this method.
|
||||
*/
|
||||
@Test
|
||||
void pendingAskInternalBypassSeesAnyDelegationsPendingQuestion() throws Exception {
|
||||
Principal creator = Principal.worker("term_creator", 1);
|
||||
String ticket = messages.sendAsync(T, "task that asks", null, creator);
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
MessageService.PendingAsk bypassed = messages.pendingAsk(T, MessageService.INTERNAL_NO_OWNER_CHECK);
|
||||
assertNotNull(bypassed, "the internal bypass must see a question on a delegation a named worker created");
|
||||
assertEquals("which config file?", bypassed.question());
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000, creator.ownerKey()));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
@@ -2037,6 +2137,118 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
// --- outstanding(): fleet_handover{open}'s list of open tickets and asks -------------------
|
||||
|
||||
@Test
|
||||
void outstandingReportsOwnedPendingTicketsWithPhases() {
|
||||
Principal lead = Principal.leader("opus", "term_lead", 1);
|
||||
String ticket1 = messages.sendAsync(T, "task one", null, lead);
|
||||
String ticket2 = messages.sendAsync(T, "task two", null, lead);
|
||||
|
||||
MessageService.Outstanding outstanding = messages.outstanding(lead.ownerKey());
|
||||
|
||||
assertEquals(2, outstanding.tickets().size());
|
||||
assertTrue(outstanding.tickets().stream().anyMatch(t ->
|
||||
ticket1.equals(t.ticket()) && t.phase() == MessageService.Phase.PENDING && T.equals(t.target())),
|
||||
"ticket1 must be reported PENDING: " + outstanding.tickets());
|
||||
assertTrue(outstanding.tickets().stream().anyMatch(t ->
|
||||
ticket2.equals(t.ticket()) && t.phase() == MessageService.Phase.PENDING && T.equals(t.target())),
|
||||
"ticket2 must be reported PENDING: " + outstanding.tickets());
|
||||
assertTrue(outstanding.asks().isEmpty(), "neither ticket has an open question");
|
||||
}
|
||||
|
||||
@Test
|
||||
void outstandingReportsAnOpenAsksTurnId() throws Exception {
|
||||
Principal lead = Principal.leader("opus", "term_lead", 1);
|
||||
String ticket = messages.sendAsync(T, "task that asks", null, lead);
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
MessageService.Outstanding outstanding = messages.outstanding(lead.ownerKey());
|
||||
assertEquals(1, outstanding.tickets().size());
|
||||
assertEquals(MessageService.Phase.ASKING, outstanding.tickets().get(0).phase());
|
||||
assertEquals(1, outstanding.asks().size());
|
||||
MessageService.OutstandingAsk openAsk = outstanding.asks().get(0);
|
||||
assertEquals(ticket, openAsk.ticket());
|
||||
assertEquals(asking.turnId(), openAsk.turnId());
|
||||
assertEquals(T, openAsk.workerSession());
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000, lead.ownerKey()));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
/**
|
||||
* Positive control: {@code leadA} must still see its own ticket and ask, so {@code leadB}
|
||||
* seeing neither is the owner filter at work and not {@code outstanding} refusing everyone.
|
||||
*/
|
||||
@Test
|
||||
void outstandingDoesNotLeakAcrossNamedLeads() throws Exception {
|
||||
Principal leadA = Principal.leader("opus", "term_a", 1);
|
||||
Principal leadB = Principal.leader("sol", "term_b", 2);
|
||||
String ticket = messages.sendAsync(T, "task that asks", null, leadA);
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||
|
||||
MessageService.Outstanding seenByA = messages.outstanding(leadA.ownerKey());
|
||||
assertEquals(1, seenByA.tickets().size(), "lead A must see its own ticket");
|
||||
assertEquals(1, seenByA.asks().size(), "lead A must see its own open ask");
|
||||
|
||||
MessageService.Outstanding seenByB = messages.outstanding(leadB.ownerKey());
|
||||
assertTrue(seenByB.tickets().isEmpty(), "lead B must not see lead A's ticket");
|
||||
assertTrue(seenByB.asks().isEmpty(), "lead B must not see lead A's open ask");
|
||||
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
() -> messages.answer(asking.turnId(), "config.yaml", 5000, leadA.ownerKey()));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting();
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
@Test
|
||||
void outstandingIsEmptyCollectionsNotNullForACallerWithNothing() {
|
||||
MessageService.Outstanding outstanding =
|
||||
messages.outstanding(Principal.leader("opus", "term_lead", 1).ownerKey());
|
||||
assertNotNull(outstanding.tickets(), "a caller with no delegations still gets a list, not null");
|
||||
assertNotNull(outstanding.asks(), "a caller with no delegations still gets a list, not null");
|
||||
assertTrue(outstanding.tickets().isEmpty());
|
||||
assertTrue(outstanding.asks().isEmpty());
|
||||
}
|
||||
|
||||
/**
|
||||
* A ticket is destroyed on a timer once it goes terminal ({@link #pruneTerminalTickets}'s TTL
|
||||
* runs from completion), so it is the one case where carrying the id forward actually matters
|
||||
* — a still-PENDING ticket is in no such danger, its worker is still running.
|
||||
*/
|
||||
@Test
|
||||
void outstandingReportsACompletedUncollectedTicketWithATerminalPhase() throws Exception {
|
||||
Principal lead = Principal.leader("opus", "term_lead", 1);
|
||||
String ticket = messages.sendAsync(T, "long task", null, lead);
|
||||
awaitWaiting();
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "async result"), "a reply resolves the async send");
|
||||
awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||
|
||||
MessageService.Outstanding outstanding = messages.outstanding(lead.ownerKey());
|
||||
assertEquals(1, outstanding.tickets().size());
|
||||
MessageService.OutstandingTicket done = outstanding.tickets().get(0);
|
||||
assertEquals(ticket, done.ticket());
|
||||
assertEquals(MessageService.Phase.DONE, done.phase());
|
||||
assertEquals(T, done.target());
|
||||
}
|
||||
|
||||
// --- fleetd #715: answer() is gated on the caller that owns the turn -----------------------
|
||||
|
||||
/**
|
||||
|
||||
@@ -207,14 +207,15 @@ class FleetAppAuthTest {
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code GET /tasks/{ticket}} refuses a worker whose terminal did not create the ticket,
|
||||
* while the creating worker and the unnamed primary both still read it. The ticket is minted
|
||||
* directly on the shared {@link MessageService}, the same way {@code MessageServiceTest}
|
||||
* drives {@link MessageService#poll(String, String)}, so this exercises only the REST poll
|
||||
* route's own handling of the ownership already recorded on the ticket.
|
||||
* {@code GET /tasks/{ticket}} refuses a worker whose terminal did not create the ticket, and
|
||||
* refuses an unnamed primary just the same: a named worker's ticket is not anyone else's to
|
||||
* read, caller rank included. The ticket is minted directly on the shared
|
||||
* {@link MessageService}, the same way {@code MessageServiceTest} drives
|
||||
* {@link MessageService#poll(String, String)}, so this exercises only the REST poll route's
|
||||
* own handling of the ownership already recorded on the ticket.
|
||||
*/
|
||||
@Test
|
||||
void restPollRefusesADifferentWorkerButAllowsTheCreatorAndTheUnnamedPrimary() throws Exception {
|
||||
void restPollRefusesADifferentWorkerAndAnUnnamedPrimaryOnANamedWorkersTicket() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Injector injector = new Injector(agents);
|
||||
@@ -241,8 +242,10 @@ class FleetAppAuthTest {
|
||||
|
||||
HttpResponse<String> primary = send(primaryApp.port(), "GET", "/tasks/" + ticket, null, null);
|
||||
assertEquals(200, primary.statusCode());
|
||||
assertFalse(primary.body().contains("forbidden"),
|
||||
"the unnamed primary must read any ticket: " + primary.body());
|
||||
assertTrue(primary.body().contains("forbidden"),
|
||||
"an unnamed primary must not read a ticket a named worker created: " + primary.body());
|
||||
assertFalse(primary.body().contains("\"reply\""),
|
||||
"a refusal must never carry reply text: " + primary.body());
|
||||
} finally {
|
||||
creatorApp.stop();
|
||||
otherWorkerApp.stop();
|
||||
@@ -250,6 +253,33 @@ class FleetAppAuthTest {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Positive control for
|
||||
* {@link #restPollRefusesADifferentWorkerAndAnUnnamedPrimaryOnANamedWorkersTicket}: without
|
||||
* this, that test's refusal would pass just as well if the route refused every caller. Here
|
||||
* the ticket's creator is itself an unnamed primary, so another unnamed primary reading it
|
||||
* over REST must still succeed.
|
||||
*/
|
||||
@Test
|
||||
void restPollAllowsAnUnnamedPrimaryItsOwnTicket() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Injector injector = new Injector(agents);
|
||||
MessageService messages = new MessageService(agents, injector, new Rendezvous());
|
||||
|
||||
Javalin primaryApp = startOnSharedService(messages, herdr, 999_999L); // no pane -> primary
|
||||
try {
|
||||
String ticket = messages.sendAsync("term_a", "long task", null, Principal.primary(FakeHerdr.WORKER_PID));
|
||||
|
||||
HttpResponse<String> own = send(primaryApp.port(), "GET", "/tasks/" + ticket, null, null);
|
||||
assertEquals(200, own.statusCode());
|
||||
assertFalse(own.body().contains("forbidden"),
|
||||
"an unnamed primary must read a ticket another unnamed primary created: " + own.body());
|
||||
} finally {
|
||||
primaryApp.stop();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* {@code POST /sessions/{id}/message} with {@code wait:false} must record the creating
|
||||
* caller's own terminal on the ticket it returns, so that caller can still poll its own
|
||||
@@ -290,9 +320,10 @@ class FleetAppAuthTest {
|
||||
|
||||
/**
|
||||
* {@code GET /sessions/{id}/status} shows a worker's pending {@code fleet_ask} question, its
|
||||
* {@code turnId} and its ticket only to the caller whose owner key created that delegation, or
|
||||
* to the unnamed primary. A different caller still sees the base status line, but none of the
|
||||
* pending-ask fields.
|
||||
* {@code turnId} and its ticket only to the caller whose owner key created that delegation. An
|
||||
* unnamed primary is held to the same rule: its owner key is {@code null}, which here does not
|
||||
* match the named worker that created the delegation, so it sees none of the pending-ask
|
||||
* fields either — the same as any other non-creating caller.
|
||||
*/
|
||||
@Test
|
||||
void restStatusGatesThePendingAskFieldsByTheDelegationsCreatorOwner() throws Exception {
|
||||
@@ -321,7 +352,8 @@ class FleetAppAuthTest {
|
||||
MessageService.TaskView asking;
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
do {
|
||||
asking = messages.poll(ticket, null);
|
||||
// the no-check overload: this is test plumbing waiting for ASKING, not the gate under test
|
||||
asking = messages.poll(ticket);
|
||||
Thread.sleep(5);
|
||||
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
|
||||
assertEquals(MessageService.Phase.ASKING, asking.phase());
|
||||
@@ -342,8 +374,11 @@ class FleetAppAuthTest {
|
||||
|
||||
JsonNode primary = mapper.readTree(
|
||||
send(primaryApp.port(), "GET", "/sessions/term_target/status", null, null).body());
|
||||
assertEquals("which config file?", primary.get("question").asText(),
|
||||
"a caller with no terminal (the unnamed primary) must see the question");
|
||||
assertEquals("idle", primary.get("status").asText(), "the base status must still be shown");
|
||||
assertFalse(primary.has("question"),
|
||||
"an unnamed primary must not see a question on a delegation a named worker created: " + primary);
|
||||
assertFalse(primary.has("turnId"), "a non-creating unnamed primary must not see the turnId: " + primary);
|
||||
assertFalse(primary.has("ticket"), "a non-creating unnamed primary must not see the ticket: " + primary);
|
||||
|
||||
// Clean up the still-open ask so the background thread does not linger past the test.
|
||||
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||
@@ -398,7 +433,8 @@ class FleetAppAuthTest {
|
||||
MessageService.TaskView asking;
|
||||
deadline = System.currentTimeMillis() + 3000;
|
||||
do {
|
||||
asking = messages.poll(ticket, null);
|
||||
// the no-check overload: this is test plumbing waiting for ASKING, not the gate under test
|
||||
asking = messages.poll(ticket);
|
||||
Thread.sleep(5);
|
||||
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
|
||||
assertEquals(MessageService.Phase.ASKING, asking.phase());
|
||||
|
||||
Reference in New Issue
Block a user