Compare commits
16 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 0d5944af63 | |||
| 65a78932c1 | |||
| 73aab3f83e | |||
| 887aca0183 | |||
| 3fd23ecafa | |||
| f379847942 | |||
| ea12107497 | |||
| 591df91de1 | |||
| 6a814176f0 | |||
| d11d1d157c | |||
| d057d56156 | |||
| b32a30fd47 | |||
| 464dbc0930 | |||
| a8cadd9150 | |||
| 57b8c0b56d | |||
| 147f50c19e |
@@ -1,10 +1,12 @@
|
||||
package dev.ltms.fleet.inject;
|
||||
|
||||
/**
|
||||
* Notified when {@link CompletionResolver} actually delivers a typed backend-error classification
|
||||
* to a waiting send (fleetd#201 / #227) — never on a race that lost. {@link CompletionResolver}
|
||||
* calls this only after {@code Rendezvous.resolveFailure} returns {@code true} for that exact
|
||||
* waiter, mirroring the win-only race rule {@link ExhaustionSink} already uses.
|
||||
* Notified when {@link CompletionResolver} has a backend-error match at the start of a pane line,
|
||||
* or both a match and its too-fast crash signature, for a waiting send (fleetd#201 / #227). A text
|
||||
* match inside ordinary pane prose can be a member's report about an error, so it fails the send
|
||||
* without notifying this sink.
|
||||
* {@link CompletionResolver} calls this only after {@code Rendezvous.resolveFailure} returns
|
||||
* {@code true} for that exact waiter, mirroring the win-only race rule {@link ExhaustionSink} uses.
|
||||
*
|
||||
* <p>The public send result is unchanged by this classification — it is still a failed send
|
||||
* ({@code Rendezvous.Kind#FAILED}); this sink is the internal seam a later stage (fleetd#201 Unit
|
||||
|
||||
@@ -387,9 +387,10 @@ public final class CompletionResolver implements TurnListener {
|
||||
// fleetd#164 (part 2) / fleetd#201: a scrape that read cleanly and produced content still
|
||||
// isn't a real reply when that content is the backend's own rejection (e.g. an HTTP 400
|
||||
// before the worker did any work). Classify it as a failure naming the member, rather than
|
||||
// handing the caller a scrape that reads like a completed answer, and — only on the
|
||||
// resolution that actually wins the race, mirroring the exhaustion sink above — notify the
|
||||
// typed backend-error sink so a later stage can act on repeated failures.
|
||||
// handing the caller a scrape that reads like a completed answer. A text match alone is not
|
||||
// enough to notify the typed backend-error sink: this assistant block can be a member's
|
||||
// normal prose about an error. A line that starts with the error match is stronger evidence;
|
||||
// the too-fast path below also has its crash signature before it records a credential failure.
|
||||
String backendError = firstMatchingLine(assistantBlock, backendErrorPatternOrFallback(target));
|
||||
if (backendError != null) {
|
||||
// Carry the whole scrape, not just the matched line. The pattern is a heuristic: a member
|
||||
@@ -401,9 +402,9 @@ public final class CompletionResolver implements TurnListener {
|
||||
if (rendezvous.resolveFailure(waiter, reason)) {
|
||||
inFlight.remove(target, turn);
|
||||
log.warn("failing send to {} via turn-stall fallback: {}", target, reason);
|
||||
// fleetd#201 Unit 1: only on the resolution that actually won the race — a late
|
||||
// duplicate must never double-count one backend failure.
|
||||
backendErrorSink.onBackendError(target, backendError, reason);
|
||||
if (startsWithBackendError(backendError, backendErrorPatternOrFallback(target))) {
|
||||
backendErrorSink.onBackendError(target, backendError, reason);
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
@@ -492,8 +493,9 @@ public final class CompletionResolver implements TurnListener {
|
||||
if (rendezvous.resolveFailure(waiter, reason)) {
|
||||
inFlight.remove(target, turn);
|
||||
log.warn("failing send to {} via turn-stall fallback from the raw scrape: {}", target, reason);
|
||||
// fleetd#201 Unit 1: only on the resolution that actually won the race.
|
||||
backendErrorSink.onBackendError(target, backendError, reason);
|
||||
if (startsWithBackendError(backendError, backendErrorPatternOrFallback(target))) {
|
||||
backendErrorSink.onBackendError(target, backendError, reason);
|
||||
}
|
||||
}
|
||||
return true;
|
||||
}
|
||||
@@ -540,10 +542,10 @@ public final class CompletionResolver implements TurnListener {
|
||||
* {@link #MIN_TURN_NANOS} — a crash signature (e.g. a backend HTTP 400 before the worker did
|
||||
* anything) that a bare {@code BUSY -> DONE} transition cannot be told apart from a genuinely
|
||||
* fast completion. Runs the same backend-error classification the normal and raw-scrape paths
|
||||
* apply, against whatever is on screen right now: a match is a typed failure that notifies
|
||||
* {@link #backendErrorSink} (only on the resolution that wins the race); a non-match stays the
|
||||
* original generic too-fast failure, naming the member and both timings, with whatever the pane
|
||||
* shows appended so the caller sees the cause, not just "it failed".
|
||||
* apply, against whatever is on screen right now: a match together with the too-fast crash
|
||||
* signature notifies {@link #backendErrorSink} (only on the resolution that wins the race). A
|
||||
* non-match stays the original generic too-fast failure, naming the member and both timings,
|
||||
* with whatever the pane shows appended so the caller sees the cause, not just "it failed".
|
||||
*/
|
||||
private void failTooFast(String target, InFlight turn, CompletableFuture<Rendezvous.Resolution> waiter,
|
||||
long elapsedNanos) {
|
||||
@@ -611,6 +613,28 @@ public final class CompletionResolver implements TurnListener {
|
||||
return null;
|
||||
}
|
||||
|
||||
/**
|
||||
* True when the error pattern begins the matched pane line, rather than appearing in prose.
|
||||
*
|
||||
* <p>Leading terminal chrome is skipped first — box-drawing characters, bullets, gutter bars and
|
||||
* spaces. #339 introduced this check with a bare {@code lookingAt}, and that rejected a genuine
|
||||
* error line rendered as {@code "| 503 Service Unavailable: ..."}: the send still failed, but the
|
||||
* credential outage was never recorded. That is the false negative #339's own invariant 3 called
|
||||
* worse than the false positive it set out to fix — measured with a throwaway probe on the
|
||||
* raw-scrape path, which is exactly the path whose comment says to expect leading chrome.
|
||||
*
|
||||
* <p>Skipping only a leading run of non-letter, non-digit characters keeps the fix's intent. A
|
||||
* member's prose ({@code "I checked the retry path. An API Error: makes it back off."}) still
|
||||
* does not match, because there the pattern sits after words, not after chrome.
|
||||
*/
|
||||
private static boolean startsWithBackendError(String line, Pattern pattern) {
|
||||
int i = 0;
|
||||
while (i < line.length() && !Character.isLetterOrDigit(line.charAt(i))) {
|
||||
i++;
|
||||
}
|
||||
return pattern.matcher(line.substring(i)).lookingAt();
|
||||
}
|
||||
|
||||
/**
|
||||
* Coverage summary for the CB-578 stage A exhausted-pattern classification, logged at startup
|
||||
* the way {@link dev.ltms.fleet.health.FleetHealthMonitor#coverage} is — so an operator can
|
||||
|
||||
@@ -145,8 +145,58 @@ public final class Injector {
|
||||
return router != null ? router.agentsFor(target) : agents;
|
||||
}
|
||||
|
||||
/** The result of trying to remove an undelivered message from the injector. */
|
||||
public enum Cancellation {
|
||||
CANCELLED,
|
||||
DELIVERED,
|
||||
NOT_DELIVERED
|
||||
}
|
||||
|
||||
/**
|
||||
* An identity handle for one queued delivery. It is the only value accepted by
|
||||
* {@link #cancel(Delivery)}, so a caller cannot cancel a different message with the same target
|
||||
* or text.
|
||||
*/
|
||||
public static final class Delivery {
|
||||
private final Pending pending;
|
||||
|
||||
private Delivery(Pending pending) {
|
||||
this.pending = pending;
|
||||
}
|
||||
|
||||
public CompletableFuture<Void> completion() {
|
||||
return pending.delivered;
|
||||
}
|
||||
}
|
||||
|
||||
/** A pending message and the future that completes when it has been delivered. */
|
||||
private record Pending(String text, TurnToken token, CompletableFuture<Void> delivered) {
|
||||
private static final class Pending {
|
||||
enum State { QUEUED, DELIVERED, NOT_DELIVERED, CANCELLED }
|
||||
|
||||
final String target;
|
||||
final String text;
|
||||
final TurnToken token;
|
||||
final CompletableFuture<Void> delivered;
|
||||
volatile State state = State.QUEUED; // written under the owning Target monitor
|
||||
|
||||
Pending(String target, String text, TurnToken token, CompletableFuture<Void> delivered) {
|
||||
this.target = target;
|
||||
this.text = text;
|
||||
this.token = token;
|
||||
this.delivered = delivered;
|
||||
}
|
||||
|
||||
String text() {
|
||||
return text;
|
||||
}
|
||||
|
||||
TurnToken token() {
|
||||
return token;
|
||||
}
|
||||
|
||||
CompletableFuture<Void> delivered() {
|
||||
return delivered;
|
||||
}
|
||||
}
|
||||
|
||||
/** Per-worker delivery state, guarded by its own monitor (single writer per worker). */
|
||||
@@ -177,15 +227,47 @@ public final class Injector {
|
||||
* <p>Uses an atomic map update so a concurrent {@link #drop} cannot slip between "find the
|
||||
* target" and "queue the message" and orphan it in a target it just removed.
|
||||
*/
|
||||
public CompletableFuture<Void> enqueue(String target, String text, TurnToken token) {
|
||||
public Delivery enqueue(String target, String text, TurnToken token) {
|
||||
CompletableFuture<Void> delivered = new CompletableFuture<>();
|
||||
Pending p = new Pending(text, token, delivered);
|
||||
Pending p = new Pending(target, text, token, delivered);
|
||||
targets.compute(target, (_, existing) -> {
|
||||
Target t = (existing != null) ? existing : new Target();
|
||||
t.add(p); // synchronized on the Target monitor — atomic with a concurrent drop
|
||||
return t;
|
||||
});
|
||||
return delivered;
|
||||
return new Delivery(p);
|
||||
}
|
||||
|
||||
/**
|
||||
* Cancel this exact queued delivery. The target monitor serializes this operation with
|
||||
* {@link #onStatus}: if delivery wins that race, this returns {@link Cancellation#DELIVERED}
|
||||
* rather than claiming the message remained queued.
|
||||
*/
|
||||
public Cancellation cancel(Delivery delivery) {
|
||||
Pending p = delivery.pending;
|
||||
Target t = targets.get(p.target);
|
||||
if (t == null) {
|
||||
return cancellationOf(p);
|
||||
}
|
||||
synchronized (t) {
|
||||
if (p.state != Pending.State.QUEUED || !t.queue.remove(p)) {
|
||||
return cancellationOf(p);
|
||||
}
|
||||
p.state = Pending.State.CANCELLED;
|
||||
if (isQuiescent(t)) {
|
||||
targets.remove(p.target, t);
|
||||
}
|
||||
return Cancellation.CANCELLED;
|
||||
}
|
||||
}
|
||||
|
||||
private static Cancellation cancellationOf(Pending p) {
|
||||
return p.state == Pending.State.DELIVERED ? Cancellation.DELIVERED : Cancellation.NOT_DELIVERED;
|
||||
}
|
||||
|
||||
private static boolean isQuiescent(Target t) {
|
||||
return t.queue.isEmpty() && !t.awaitingPickup && !t.awaitingCompletion
|
||||
&& !t.postTurnPending && !t.awaitingPostTurnPickup && !t.postTurnObserved;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -274,6 +356,7 @@ public final class Injector {
|
||||
try {
|
||||
agentsFor(target).send(target, p.text());
|
||||
t.queue.poll();
|
||||
p.state = Pending.State.DELIVERED;
|
||||
t.awaitingPickup = true;
|
||||
t.awaitingCompletion = true;
|
||||
t.turnObserved = false;
|
||||
@@ -283,6 +366,7 @@ public final class Injector {
|
||||
// Delivery failed at herdr; drop the poisoned message and surface it
|
||||
// rather than blocking the queue behind it.
|
||||
t.queue.poll();
|
||||
p.state = Pending.State.NOT_DELIVERED;
|
||||
sent = p;
|
||||
sendError = e;
|
||||
}
|
||||
@@ -293,6 +377,9 @@ public final class Injector {
|
||||
// fail every queued message and release the target (CB-114) instead of
|
||||
// polling it indefinitely with the caller's future never completing.
|
||||
notReady = new ArrayList<>(t.queue);
|
||||
for (Pending pending : notReady) {
|
||||
pending.state = Pending.State.NOT_DELIVERED;
|
||||
}
|
||||
log.warn("readiness grace for {} expired after {} polls ({}s): target never "
|
||||
+ "became deliverable, so failing {} queued message(s) that never "
|
||||
+ "reached its pane",
|
||||
@@ -340,8 +427,7 @@ public final class Injector {
|
||||
|
||||
// Reclaim the entry once the worker is fully quiescent (nothing queued, no pickup or
|
||||
// completion awaited), so the map cannot grow without bound across short-lived workers.
|
||||
if (t.queue.isEmpty() && !t.awaitingPickup && !t.awaitingCompletion
|
||||
&& !t.postTurnPending && !t.awaitingPostTurnPickup && !t.postTurnObserved) {
|
||||
if (isQuiescent(t)) {
|
||||
targets.remove(target, t);
|
||||
}
|
||||
}
|
||||
@@ -432,6 +518,9 @@ public final class Injector {
|
||||
boolean hadDeliveredTurn;
|
||||
synchronized (t) {
|
||||
pending = new ArrayList<>(t.queue);
|
||||
for (Pending p : pending) {
|
||||
p.state = Pending.State.NOT_DELIVERED;
|
||||
}
|
||||
t.queue.clear();
|
||||
hadDeliveredTurn = t.awaitingCompletion;
|
||||
t.awaitingCompletion = false;
|
||||
|
||||
@@ -937,14 +937,20 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
*/
|
||||
@Override
|
||||
public void stop(String idOrPane) {
|
||||
// Teardown knows only the pane, not which profile spawned it. Attempt tab cleanup when any
|
||||
// profile uses tab placement (so the bridge may have created a dedicated peer tab); the
|
||||
// single-occupant check below is what actually protects the user's shared tabs.
|
||||
// Teardown knows only the pane, not which profile spawned it — and, when this call is
|
||||
// routed here through CompositePeerLauncher's single-daemon stop() shortcut (fleetd #342),
|
||||
// not even which adapter's config actually governed the spawn: the shortcut can hand the
|
||||
// pane to a delegate that never spawned it, whose own profiles say nothing about how THIS
|
||||
// pane was placed. So the decision to look for a tab to clean up is made from the pane's
|
||||
// actual state, not from this delegate's static profile config: resolve the tab
|
||||
// unconditionally and let {@link WorkspaceControl#locatePane} tolerate "not found" (it
|
||||
// returns null rather than throwing); the single-occupant check below is what actually
|
||||
// protects the user's shared tabs, exactly as it always has.
|
||||
String paneId = paneByAgentId.remove(idOrPane);
|
||||
if (paneId == null) {
|
||||
paneId = idOrPane; // raw-pane fallback (reap, gate timeout, pane-addressed callers)
|
||||
}
|
||||
WorkspaceControl.PaneLocation loc = usesTabPlacement() ? spaces.locatePane(paneId) : null;
|
||||
WorkspaceControl.PaneLocation loc = spaces.locatePane(paneId);
|
||||
try {
|
||||
agents.close(paneId);
|
||||
} catch (HerdrException e) {
|
||||
@@ -985,11 +991,6 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
}
|
||||
}
|
||||
|
||||
/** Whether any configured profile places peers in their own tab (so tabs may need cleanup). */
|
||||
private boolean usesTabPlacement() {
|
||||
return profiles.values().stream().anyMatch(FleetConfig.Profile::tabPlacement);
|
||||
}
|
||||
|
||||
/** True when a herdr error means the target is already gone (safe to treat as done). */
|
||||
private static boolean isAlreadyGone(HerdrException e) {
|
||||
return e.code() != null && e.code().endsWith("_not_found");
|
||||
@@ -1665,30 +1666,59 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
|
||||
/**
|
||||
* Guards the {@code effectiveAllowed == null} branch of {@link #logCredentialGap} — the
|
||||
* genuinely-unprotected report (deny-by-default, and the allow-list non-zsh fallback) — to one
|
||||
* WARN per launcher instance, not one per spawn.
|
||||
* genuinely-unprotected report (deny-by-default, and the allow-list non-zsh fallback) — AND
|
||||
* the {@code effectiveAllowed != null} / {@code keptByDerivedList} branch, the allow-list case
|
||||
* where a name is in the gap but the derived allow-list keeps it anyway. Both branches log the
|
||||
* same severity (WARN) about the same fact — a name genuinely reaching a member pane
|
||||
* unprotected — so they share this one guard, keyed per NAME rather than per launcher instance:
|
||||
* each credential-shaped name that is ever reported unprotected gets exactly one WARN, however
|
||||
* many spawns see it and whichever of the two branches first reports it.
|
||||
*
|
||||
* <p>fleetd #341: {@code memberCredentials} is a live, re-read-per-spawn supplier, so the
|
||||
* policy — and so the gap's actual member names — can change between two spawns on the same
|
||||
* launcher. Before this fix the guard was a single {@code AtomicBoolean} tripped by either
|
||||
* branch: spawn 1 could warn about name A and trip the flag, and a later spawn's gap containing
|
||||
* a different name B would never be reported, even though B is just as unprotected as A was.
|
||||
* {@code AtomicBoolean} could not express "once per distinct name" at all — only "once, ever,
|
||||
* for whichever name got there first" — so this is a {@code Set<String>} guard instead, the
|
||||
* same shape {@link OpenCodeLauncher#modelCheckSkippedWarned} already uses for its own
|
||||
* once-per-distinct-thing WARN. {@link #add}'s return value (true only the first time a name is
|
||||
* added) is what turns "log the whole gap" into "log only the names never warned about before".
|
||||
*
|
||||
* <p>Bounded by construction: every name added here first passed {@link
|
||||
* #CREDENTIAL_SHAPED_NAME}'s filter over {@link #hostEnvNames}, i.e. it is an actual
|
||||
* environment variable name from the daemon's own process — a small, OS-bounded set (the host
|
||||
* environment has, in practice, tens to a few hundred entries), not an attacker- or
|
||||
* request-controlled input. So this set cannot grow past "however many distinct credential-
|
||||
* shaped names this host's environment has ever held across this launcher's lifetime," which is
|
||||
* effectively fixed for the life of one daemon process — no separate cap is needed.
|
||||
*
|
||||
* <p>CB-633 follow-up (#192): kept SEPARATE from {@link #allowListGapLogged} on purpose.
|
||||
* {@code memberCredentials} is a live, re-read-per-spawn supplier, so the policy can change
|
||||
* between two spawns on the same launcher. A single shared flag would let a harmless allow-list
|
||||
* INFO on spawn 1 permanently suppress the real deny-by-default WARN a later spawn deserves —
|
||||
* the report that matters most getting hidden by the report that doesn't. Two flags mean each
|
||||
* report kind fires exactly once, independent of what the other kind already logged.
|
||||
* report kind (WARN vs. INFO) fires independently of what the other kind already logged; within
|
||||
* the WARN kind itself, the set above further separates by name, for the same reason.
|
||||
*/
|
||||
private final AtomicBoolean unprotectedGapLogged = new AtomicBoolean();
|
||||
private final Set<String> unprotectedGapNamesWarned = ConcurrentHashMap.newKeySet();
|
||||
|
||||
/**
|
||||
* Guards the {@code effectiveAllowed != null} branch of {@link #logCredentialGap} — the
|
||||
* allow-list-scrub-covered report — to one INFO per launcher instance. See {@link
|
||||
* #unprotectedGapLogged}'s javadoc for why this is a separate flag rather than a shared one.
|
||||
* #unprotectedGapNamesWarned}'s javadoc for why this is a separate flag rather than a shared
|
||||
* one; unlike that guard it stays a per-instance {@code AtomicBoolean}, not a per-name set —
|
||||
* fleetd #341 fixed the WARN-vs-WARN suppression, not this INFO's own one-shot shape, which was
|
||||
* not reported as broken and is out of that ticket's scope.
|
||||
*/
|
||||
private final AtomicBoolean allowListGapLogged = new AtomicBoolean();
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 2: guards {@link #warnUnknownMemberEnvironment} to one WARN per launcher
|
||||
* instance, not one per spawn — the same one-per-instance shape as {@link #unprotectedGapLogged}
|
||||
* and {@link #allowListGapLogged}, kept as its own flag for the same reason those two are split:
|
||||
* this mode is orthogonal to which of the other two branches would otherwise have fired.
|
||||
* instance, not one per spawn — the same one-shot shape {@link #unprotectedGapNamesWarned} and
|
||||
* {@link #allowListGapLogged} guard their own branches with, kept as its own flag for the same
|
||||
* reason those two are split: this mode is orthogonal to which of the other two branches would
|
||||
* otherwise have fired.
|
||||
*/
|
||||
private final AtomicBoolean unknownMemberEnvironmentWarned = new AtomicBoolean();
|
||||
|
||||
@@ -1816,14 +1846,21 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
List<String> blankedByScrub = gap.stream()
|
||||
.filter(name -> !MemberEnvAllowList.keeps(effectiveAllowed, name))
|
||||
.toList();
|
||||
if (!keptByDerivedList.isEmpty() && unprotectedGapLogged.compareAndSet(false, true)) {
|
||||
// fleetd #341: filter to names this guard has never warned about before — not just
|
||||
// "isEmpty" on the whole branch — so a name this spawn's gap shares with an EARLIER
|
||||
// spawn's (already-warned) gap does not re-print, while a name unique to THIS gap still
|
||||
// does, whichever of the two WARN branches reported it first.
|
||||
List<String> newlyUnprotected = keptByDerivedList.stream()
|
||||
.filter(unprotectedGapNamesWarned::add)
|
||||
.toList();
|
||||
if (!newlyUnprotected.isEmpty()) {
|
||||
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
|
||||
+ "known: nor allow: — the derived allow-list keeps them anyway (a profile's "
|
||||
+ "gitTokenEnv/gitHostEnv/tokenEnv/env: names one, or this spawn injects it), "
|
||||
+ "so every member pane inherits them UNBLOCKED — {}. Add each to "
|
||||
+ "memberCredentials.known (or .allow if a member legitimately needs it), or "
|
||||
+ "remove it from whatever profile setting derives it in.",
|
||||
keptByDerivedList.size(), keptByDerivedList);
|
||||
newlyUnprotected.size(), newlyUnprotected);
|
||||
}
|
||||
if (!blankedByScrub.isEmpty() && allowListGapLogged.compareAndSet(false, true)) {
|
||||
log.info("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
|
||||
@@ -1836,12 +1873,18 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
||||
|
||||
/** The deny-by-default (and allow-list non-zsh fallback) WARN — unchanged byte-for-byte by #192. */
|
||||
private void warnGapUnprotected(List<String> gap) {
|
||||
if (unprotectedGapLogged.compareAndSet(false, true)) {
|
||||
// fleetd #341: same "only the names never warned before" filter as the sibling branch in
|
||||
// logCredentialGap above — see unprotectedGapNamesWarned's javadoc. Both branches share
|
||||
// this one guard because both report the exact same fact (a name reaching a member pane
|
||||
// unprotected) at the exact same severity; keying it by name is what lets a later spawn's
|
||||
// DIFFERENT name still get its own WARN after an earlier spawn's already fired.
|
||||
List<String> newlyUnprotected = gap.stream().filter(unprotectedGapNamesWarned::add).toList();
|
||||
if (!newlyUnprotected.isEmpty()) {
|
||||
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
|
||||
+ "known: nor allow: — every member pane inherits them UNBLOCKED — {}. "
|
||||
+ "Add each to memberCredentials.known (blocked by default) or .allow "
|
||||
+ "(if a member legitimately needs it).",
|
||||
gap.size(), gap);
|
||||
newlyUnprotected.size(), newlyUnprotected);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -709,27 +709,47 @@ public final class MessageService {
|
||||
boolean isRecovery = task == recoveryTask && recovered != null;
|
||||
Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
|
||||
String turnId = task.turnId;
|
||||
if (task.future.complete(outcome)) {
|
||||
if (outcome.outcome() == Outcome.WORKER_FAILED) {
|
||||
asyncFailed = true;
|
||||
// fleetd #335: completing THIS task's future must not depend on any other task's
|
||||
// cleanup succeeding — every task in `matching` is owed its own outcome regardless of
|
||||
// what happens below, so decide and record that before doing anything that can throw.
|
||||
boolean completedHere = task.future.complete(outcome);
|
||||
if (completedHere && outcome.outcome() == Outcome.WORKER_FAILED) {
|
||||
asyncFailed = true;
|
||||
}
|
||||
try {
|
||||
if (abandonCleanupHookForTest != null) {
|
||||
// Test-only (fleetd #335 site 1): see the field's own javadoc.
|
||||
abandonCleanupHookForTest.run();
|
||||
}
|
||||
if (turnId != null) {
|
||||
// #275: whether this task was swept out of ASKING or was already answered and
|
||||
// only waiting on its resumed turn's real reply (#137), nothing will ever
|
||||
// complete this turnId now — drop it from this class's own bookkeeping AND the
|
||||
// reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn
|
||||
// that is actually done, and a late answer() sees it as lapsed rather than
|
||||
// resolving a question nothing is listening for any more.
|
||||
asyncTasksByTurn.remove(turnId, task);
|
||||
rendezvous.closeAsk(turnId);
|
||||
if (completedHere) {
|
||||
if (turnId != null) {
|
||||
// #275: whether this task was swept out of ASKING or was already answered and
|
||||
// only waiting on its resumed turn's real reply (#137), nothing will ever
|
||||
// complete this turnId now — drop it from this class's own bookkeeping AND the
|
||||
// reverse-rendezvous itself, so hasAsyncQuestion(target) stops reporting a turn
|
||||
// that is actually done, and a late answer() sees it as lapsed rather than
|
||||
// resolving a question nothing is listening for any more.
|
||||
asyncTasksByTurn.remove(turnId, task);
|
||||
rendezvous.closeAsk(turnId);
|
||||
}
|
||||
} else if (isRecovery) {
|
||||
// The recovered reply was already drained out of the inbox, but this task resolved
|
||||
// through another path (e.g. a concurrent reply() or a second abandon() racing this
|
||||
// one) between us choosing it and completing it here. Put the reply back rather than
|
||||
// lose it silently — it may still belong to some other still-open task, or the next
|
||||
// caller that drains this target's inbox.
|
||||
inbox.publish(target, UUID.randomUUID().toString(), recovered.text());
|
||||
}
|
||||
} else if (isRecovery) {
|
||||
// The recovered reply was already drained out of the inbox, but this task resolved
|
||||
// through another path (e.g. a concurrent reply() or a second abandon() racing this
|
||||
// one) between us choosing it and completing it here. Put the reply back rather than
|
||||
// lose it silently — it may still belong to some other still-open task, or the next
|
||||
// caller that drains this target's inbox.
|
||||
inbox.publish(target, UUID.randomUUID().toString(), recovered.text());
|
||||
} catch (RuntimeException e) {
|
||||
// fleetd #335: inbox.publish reaches a broker (AmqpReplyInbox throws
|
||||
// IllegalStateException on an unroutable/unconfirmed/interrupted publish) and this
|
||||
// loop has no other teardown path — a caller on the release path, or the health
|
||||
// monitor's GONE/NEVER_READY sweep. Losing this exception uncaught would abort the
|
||||
// loop and leave every task still to come in `matching` PENDING forever (fleetd
|
||||
// #335 site 1). completedHere is already recorded above, so only this task's
|
||||
// best-effort bookkeeping is lost — log it and let the loop reach the rest.
|
||||
log.error("abandon: per-task cleanup failed for ticket {} (target {}, turnId {})",
|
||||
task.ticket, target, turnId, e);
|
||||
}
|
||||
}
|
||||
if (failed) {
|
||||
@@ -863,16 +883,26 @@ public final class MessageService {
|
||||
if (onAccepted != null) {
|
||||
onAccepted.run();
|
||||
}
|
||||
CompletableFuture<Void> delivered = injector.enqueue(target, content, token);
|
||||
Injector.Delivery delivery = injector.enqueue(target, content, token);
|
||||
try {
|
||||
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
|
||||
return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId()));
|
||||
} catch (TimeoutException e) {
|
||||
boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally();
|
||||
boolean wasDelivered = delivery.completion().isDone()
|
||||
&& !delivery.completion().isCompletedExceptionally();
|
||||
if (!wasDelivered) {
|
||||
if (timeoutCancellationRaceHookForTest != null) {
|
||||
// Test-only (fleetd #345): see the field's own javadoc.
|
||||
timeoutCancellationRaceHookForTest.run();
|
||||
}
|
||||
// The target monitor makes cancellation atomic with onStatus picking this
|
||||
// Pending up. If pickup won, report TIMED_OUT_WORKING because the text landed.
|
||||
wasDelivered = injector.cancel(delivery) == Injector.Cancellation.DELIVERED;
|
||||
}
|
||||
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
|
||||
if (!wasDelivered) {
|
||||
// CB-640: still sitting in the injector's queue, waiting for the member to
|
||||
// go idle — record the fact for fleet health (see queuedDeliveries).
|
||||
// CB-640: record that delivery did not happen for fleet health (see
|
||||
// queuedDeliveries). The exact Pending was cancelled, so it cannot arrive later.
|
||||
queuedDeliveries.put(target, Boolean.TRUE);
|
||||
}
|
||||
return recorded(new Reply(
|
||||
@@ -935,13 +965,41 @@ public final class MessageService {
|
||||
return new AskResult(AskOutcome.ANSWERED, answer);
|
||||
} catch (TimeoutException e) {
|
||||
log.debug("fleet_ask from {} went unanswered in {}ms", workerSession, timeoutMillis);
|
||||
// fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it out of
|
||||
// asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and stays (it is
|
||||
// what keeps the target from staying BUSY forever), but it would otherwise also erase
|
||||
// askAnsweredAsyncTasks' only signal that the worker's eventual real fleet_reply still
|
||||
// belongs to this task, stranding it in the inbox with a false "never replied" verdict.
|
||||
markAskTimedOut(ticket.turnId());
|
||||
clearAsyncQuestion(ticket.turnId(), true);
|
||||
// Only the fresh owner tears down the shared turn (mirrors the finally block below).
|
||||
// A duplicate's own timeoutMillis says nothing about whether the SHARED ask is actually
|
||||
// done — it must leave the close/forget bookkeeping to the fresh owner, exactly as it
|
||||
// already leaves closeAsk to it.
|
||||
if (ticket.fresh()) {
|
||||
// fleetd #334: close the ask turn BEFORE forgetting this task's turnId mapping below.
|
||||
// Before this fix the order was reversed — the mapping was forgotten here first, and
|
||||
// rendezvous.closeAsk only ran afterward, in the shared finally. A primary's answer()
|
||||
// call racing this exact timeout could then find rendezvous.askSession(turnId) still
|
||||
// non-null (the ask still "answerable") after the Task mapping was already gone:
|
||||
// answer()'s own asyncTasksByTurn lookup returned null, its task != null guard skipped
|
||||
// the completion, and the async ticket sat at PENDING forever even though answer()
|
||||
// itself reported the worker's real reply. Closing here first removes that window:
|
||||
// any answer() call that still observes askSession(turnId) != null is necessarily
|
||||
// racing a point BEFORE the forgetting below runs (both happen on this one thread, in
|
||||
// this order, with nothing that yields in between), so the Task mapping is still there
|
||||
// for it to find; any call that observes askSession(turnId) == null now correctly
|
||||
// bails out STALE_TURN (see answer()'s own top check) before ever reaching
|
||||
// asyncTasksByTurn. rendezvous.closeAsk is idempotent — a no-op once the turn is
|
||||
// already removed, see its own javadoc — so the shared finally below re-running it
|
||||
// for this same fresh call is harmless.
|
||||
rendezvous.closeAsk(ticket.turnId());
|
||||
if (askTimeoutRaceHookForTest != null) {
|
||||
// Test-only (fleetd #334): see the field's own javadoc.
|
||||
askTimeoutRaceHookForTest.run();
|
||||
}
|
||||
// fleetd #307: mark the task BEFORE clearAsyncQuestion(forgetTurn=true) below drops it
|
||||
// out of asyncTasksByTurn and nulls its turnId — that forgetting is deliberate and
|
||||
// stays (it is what keeps the target from staying BUSY forever), but it would
|
||||
// otherwise also erase askAnsweredAsyncTasks' only signal that the worker's eventual
|
||||
// real fleet_reply still belongs to this task, stranding it in the inbox with a false
|
||||
// "never replied" verdict.
|
||||
markAskTimedOut(ticket.turnId());
|
||||
clearAsyncQuestion(ticket.turnId(), true);
|
||||
}
|
||||
return new AskResult(AskOutcome.TIMED_OUT, null);
|
||||
} catch (ExecutionException e) {
|
||||
Throwable cause = e.getCause();
|
||||
@@ -1040,19 +1098,28 @@ public final class MessageService {
|
||||
// the chained ask deliberately left open.
|
||||
//
|
||||
// A null task is NOT only "this was never an async ticket". That reading was in this
|
||||
// comment when #329 merged and it is wrong. A genuine async ticket also lands here
|
||||
// with task == null, because ask()'s timeout path runs clearAsyncQuestion(turnId,
|
||||
// true) — which drops the asyncTasksByTurn entry — in its catch block, while
|
||||
// rendezvous.closeAsk(turnId) runs later, in its finally. Between those two the ask
|
||||
// is still answerable but the map entry is already gone, so the lookup at :991
|
||||
// returns null and this ticket is never completed. Measured on 2026-09-04: a probe
|
||||
// firing only that first half before answer() runs printed
|
||||
// "answer=REPLIED phase=PENDING reply=null" — the same stranded ticket #329 set out
|
||||
// to fix, one step earlier in the same race. The probe used forgetTurnForTest, which
|
||||
// omits ask()'s markAskTimedOut; that cannot change the outcome, because askTimedOut
|
||||
// is read only by askAnsweredAsyncTasks, and reply() never reaches it while this
|
||||
// method's own waiter is live. So #329 narrows this window rather than closing it.
|
||||
// Open as fleetd #334 — do not read this guard as complete.
|
||||
// comment when #329 merged and it is wrong; it is still not the whole story after
|
||||
// #334. A genuine async ticket can still land here with task == null — a blocking
|
||||
// (wait:true) send's fleet_ask never has a Task at all, so that case is expected and
|
||||
// fine. What #334 fixed was a SECOND, unintended way to get here with task == null:
|
||||
// ask()'s timeout path used to run clearAsyncQuestion(turnId, true) — which drops the
|
||||
// asyncTasksByTurn entry — in its catch block, while rendezvous.closeAsk(turnId) ran
|
||||
// later, in its finally. Between those two the ask was still answerable but the map
|
||||
// entry was already gone, so the lookup at :1053 returned null and this ticket was
|
||||
// never completed. Measured on 2026-09-04: a probe firing only that first half before
|
||||
// answer() ran printed "answer=REPLIED phase=PENDING reply=null" — the same stranded
|
||||
// ticket #329 set out to fix, one step earlier in the same race; the probe used
|
||||
// forgetTurnForTest, which omits ask()'s markAskTimedOut, and that omission does not
|
||||
// change the outcome, because askTimedOut is read only by askAnsweredAsyncTasks, and
|
||||
// reply() never reaches it while this method's own waiter is live. #334's fix
|
||||
// reorders ask()'s timeout catch to run closeAsk before the forgetting (see the
|
||||
// fresh-owner block there), which removes this path entirely rather than narrowing
|
||||
// it further: once closeAsk has run, rendezvous.askSession(turnId) is null and
|
||||
// answer() returns STALE_TURN from its own top check, before it ever reaches this
|
||||
// lookup — see aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket in
|
||||
// MessageServiceTest, which pins the exact window with askTimeoutRaceHookForTest.
|
||||
// So by the time this line runs, task == null means only the ordinary blocking-send
|
||||
// case (or #329's own already-fixed race elsewhere) — not this one.
|
||||
if (result.outcome() != Outcome.QUESTION && task != null) {
|
||||
finishAsyncTask(task, result);
|
||||
}
|
||||
@@ -1112,7 +1179,18 @@ public final class MessageService {
|
||||
// class javadoc on sendAsync/CB-107.
|
||||
task.future.whenComplete((reply, ex) -> {
|
||||
boolean failed = ex != null || reply == null || !reply.completed();
|
||||
pushLoop.onTicketTerminal(ticket, target, failed);
|
||||
try {
|
||||
pushLoop.onTicketTerminal(ticket, target, failed);
|
||||
} catch (Throwable t) {
|
||||
// fleetd #335 (site 2): this stage's own CompletableFuture is discarded, so an
|
||||
// uncaught throw here (e.g. a RejectedExecutionException from
|
||||
// ReplyPushLoop's scheduler, already shut down while this in-flight send's
|
||||
// whenComplete fires during the daemon's own shutdown sequence — messages.close()
|
||||
// only stops accepting NEW async work, it does not cancel a delivery already
|
||||
// running) vanishes with no log line and no metric, and the push loop never learns
|
||||
// the ticket went terminal — the exact thing this hook exists to tell it.
|
||||
log.error("push loop failed to learn ticket {} (target {}) went terminal", ticket, target, t);
|
||||
}
|
||||
});
|
||||
}
|
||||
asyncExecutor.submit(() -> {
|
||||
@@ -1364,6 +1442,24 @@ public final class MessageService {
|
||||
this.afterFinishAsyncTaskCompleteHookForTest = hook;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null in production; test seam for fleetd #345 — invoked in {@link #send}'s timeout path after
|
||||
* {@link Injector.Delivery#completion()} reports incomplete and before {@link Injector#cancel}
|
||||
* takes the target monitor. A test installs this to make {@code onStatus} pick the exact queued
|
||||
* delivery up in that window, so {@code cancel} returns {@link Injector.Cancellation#DELIVERED}.
|
||||
* This deterministically covers the caller's need to use that result rather than relying on a
|
||||
* timing-sensitive real race.
|
||||
*/
|
||||
private volatile Runnable timeoutCancellationRaceHookForTest;
|
||||
|
||||
/**
|
||||
* Test-only (fleetd #345): install {@link #timeoutCancellationRaceHookForTest}. Package-private
|
||||
* so the test, in the same package, can reach it without widening any production API.
|
||||
*/
|
||||
void setTimeoutCancellationRaceHookForTest(Runnable hook) {
|
||||
this.timeoutCancellationRaceHookForTest = hook;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null in production; test seam for fleetd #329 (F3) — invoked from {@link #reply} right after
|
||||
* the single local read of {@code orphan.turnId} passes its null-check and before that (now-local)
|
||||
@@ -1393,6 +1489,56 @@ public final class MessageService {
|
||||
clearAsyncQuestion(turnId, true);
|
||||
}
|
||||
|
||||
/**
|
||||
* Null in production; test seam for fleetd #335 (site 1) — invoked from {@link #abandon(String,
|
||||
* String, boolean)}'s per-task loop, once per task, right before that task's own cleanup
|
||||
* (turnId bookkeeping, or the stranded-reply put-back) runs. A test installs this to inject a
|
||||
* throw at that exact point deterministically.
|
||||
*
|
||||
* <p>The one production call there that can really throw is {@code inbox.publish} in the
|
||||
* put-back branch — {@link AmqpReplyInbox#publish} reaches a broker and throws {@link
|
||||
* IllegalStateException} on an unroutable, unconfirmed, or interrupted publish — but reaching
|
||||
* that branch requires a second completion of the very task {@code abandon} is about to
|
||||
* complete to win the race first (see the branch's own comment), and the #137 follow-up
|
||||
* investigation above already found the combination this needs (a stranded reply coinciding
|
||||
* with an open matching task) unreachable through the public API, not merely hard to time.
|
||||
* This hook reproduces the resulting shape — a per-task cleanup throw — directly, the same
|
||||
* technique {@link #finishAsyncTaskRaceHook} and {@link
|
||||
* #afterFinishAsyncTaskCompleteHookForTest} already use for their own hard-to-time races.
|
||||
*/
|
||||
private volatile Runnable abandonCleanupHookForTest;
|
||||
|
||||
/**
|
||||
* Test-only (fleetd #335, site 1): install {@link #abandonCleanupHookForTest}. Package-private
|
||||
* so the test, in the same package, can reach it without widening any production API.
|
||||
*/
|
||||
void setAbandonCleanupHookForTest(Runnable hook) {
|
||||
this.abandonCleanupHookForTest = hook;
|
||||
}
|
||||
|
||||
/**
|
||||
* Null in production; test seam for fleetd #334 — invoked from {@link #ask}'s {@code
|
||||
* TimeoutException} catch, only for the fresh owner, right after {@code rendezvous.closeAsk}
|
||||
* has run and before {@link #markAskTimedOut} / {@link #clearAsyncQuestion} forget this task's
|
||||
* turnId mapping. A test installs this to call {@link #answer} for the very same {@code turnId}
|
||||
* synchronously from inside that exact window, deterministically reproducing the race a real
|
||||
* concurrent {@code answer()} call could otherwise only win by timing luck: with the ask already
|
||||
* closed, that call must see {@code rendezvous.askSession(turnId) == null} and return {@link
|
||||
* Outcome#STALE_TURN} immediately, never reaching {@code asyncTasksByTurn} at all — proving the
|
||||
* window fleetd #334 describes (mapping forgotten while the ask was still "answerable") is
|
||||
* closed, rather than merely narrowed the way fleetd #329 narrowed the sibling race in {@link
|
||||
* #finishAsyncTask}.
|
||||
*/
|
||||
private volatile Runnable askTimeoutRaceHookForTest;
|
||||
|
||||
/**
|
||||
* Test-only (fleetd #334): install {@link #askTimeoutRaceHookForTest}. Package-private so the
|
||||
* test, in the same package, can reach it without widening any production API.
|
||||
*/
|
||||
void setAskTimeoutRaceHookForTest(Runnable hook) {
|
||||
this.askTimeoutRaceHookForTest = hook;
|
||||
}
|
||||
|
||||
/** A new send must not open a waiter while an async ticket owns this worker's paused turn. */
|
||||
private boolean hasAsyncQuestion(String target) {
|
||||
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
|
||||
|
||||
@@ -804,7 +804,28 @@ class CompletionResolverTest {
|
||||
// --- fleetd#201 Unit 1: target-keyed backend-error pattern + typed sink ----------------------
|
||||
|
||||
@Test
|
||||
void aConfiguredBackendErrorPatternClassifiesAMatchAsAFailureAndNotifiesTheSinkOnce() {
|
||||
void aNormalMemberReportMentioningTheFallbackErrorPatternFailsButDoesNotNotifyTheSink() {
|
||||
String block = "⏺ I checked the retry path. An API Error: makes it back off.\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
java.util.List<String> notified = new java.util.ArrayList<>();
|
||||
BackendErrorSink sink = (target, matchedLine, reason) -> notified.add(target + ": " + matchedLine);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), BackendErrorPatternLookup.legacy(), sink);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
|
||||
"a scrape mentioning the pattern must still fail the send");
|
||||
assertTrue(waiter.getNow(null).text().contains("I checked the retry path. An API Error: makes it back off."),
|
||||
"the failure must keep the whole pane tail");
|
||||
assertTrue(notified.isEmpty(),
|
||||
"a normal report mentioning the fallback pattern must not record a credential failure");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aConfiguredBackendErrorPatternAtTheStartOfALineClassifiesAMatchAndNotifiesTheSinkOnce() {
|
||||
String block = "⏺ 503 Service Unavailable: upstream credential rejected\n❯ ";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
@@ -924,7 +945,7 @@ class CompletionResolverTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void aConfiguredPatternAlsoClassifiesTheRawScrapeFallbackAndNotifiesTheSink() {
|
||||
void aConfiguredPatternAtTheStartOfALineAlsoClassifiesTheRawScrapeFallbackAndNotifiesTheSink() {
|
||||
// No ⏺ marker and leading TUI chrome ⇒ lastAssistantBlock() yields "", so classification must
|
||||
// fall back to the raw scrape (fleetd#211) — and it must use the configured pattern too.
|
||||
String block = """
|
||||
@@ -949,6 +970,40 @@ class CompletionResolverTest {
|
||||
assertTrue(notified.get(0).contains("503 Service Unavailable"), notified.get(0));
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #339 follow-up: a genuine backend error rendered behind terminal chrome must still
|
||||
* record the credential outage. #339 added a start-of-line check to stop a member's own prose
|
||||
* being counted as an outage, and a bare {@code lookingAt} also rejected this — the send failed
|
||||
* but the sink never fired. #339's invariant 3 named that direction as the worse one: a real
|
||||
* outage going unrecorded leaves the fleet spawning into a dead credential.
|
||||
*
|
||||
* <p>The raw-scrape path is where this matters, because its own comment says to expect leading
|
||||
* TUI chrome there.
|
||||
*/
|
||||
@Test
|
||||
void aRealErrorBehindTerminalChromeStillNotifiesTheSink() {
|
||||
String block = """
|
||||
╭──────────────────────────────────────╮
|
||||
│ 503 Service Unavailable: upstream credential rejected
|
||||
""";
|
||||
FakeHerdr herdr = new FakeHerdr().readText(block);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
BackendErrorPatternLookup patterns = target -> Pattern.compile("(?i)503 Service Unavailable");
|
||||
java.util.List<String> notified = new java.util.ArrayList<>();
|
||||
BackendErrorSink sink = (target, matchedLine, reason) -> notified.add(target + ": " + matchedLine);
|
||||
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(), patterns, sink);
|
||||
|
||||
var waiter = rendezvous.open("term_a");
|
||||
resolver.resolve("term_a", new CompletionResolver.InFlight(waiter, null));
|
||||
|
||||
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind(),
|
||||
"a real backend error must still fail the send");
|
||||
assertEquals(1, notified.size(),
|
||||
"a real error line behind box chrome is still a real outage — it must reach the sink, "
|
||||
+ "or the fleet keeps spawning into a dead credential");
|
||||
}
|
||||
|
||||
// --- fleetd#201 Unit 1: classification inside the fleetd#164 MIN_TURN_NANOS floor -------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -48,7 +48,7 @@ class InjectorTest {
|
||||
|
||||
@Test
|
||||
void deliversWhenIdle() {
|
||||
CompletableFuture<Void> f = injector.enqueue(T, "hello", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> f = injector.enqueue(T, "hello", TestTurnTokens.inert(T)).completion();
|
||||
assertFalse(f.isDone(), "not delivered until an injectable status arrives");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
assertTrue(f.isDone());
|
||||
@@ -170,6 +170,32 @@ class InjectorTest {
|
||||
assertEquals(List.of("a", "b", "c"), sent());
|
||||
}
|
||||
|
||||
@Test
|
||||
void cancellingTheMiddleDeliveryKeepsTheFollowingDeliveryReachable() {
|
||||
Injector.Delivery first = injector.enqueue(T, "same text", TestTurnTokens.inert(T));
|
||||
Injector.Delivery cancelled = injector.enqueue(T, "same text", TestTurnTokens.inert(T));
|
||||
injector.enqueue(T, "after cancelled", TestTurnTokens.inert(T));
|
||||
|
||||
assertEquals(Injector.Cancellation.CANCELLED, injector.cancel(cancelled));
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
assertEquals(List.of("same text", "after cancelled"), sent(),
|
||||
"cancellation must match the exact Delivery and preserve the remaining FIFO queue");
|
||||
assertTrue(first.completion().isDone());
|
||||
}
|
||||
|
||||
@Test
|
||||
void cancellationReportsDeliveredWhenPickupWonTheRace() {
|
||||
Injector.Delivery delivery = injector.enqueue(T, "already sent", TestTurnTokens.inert(T));
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
assertEquals(Injector.Cancellation.DELIVERED, injector.cancel(delivery),
|
||||
"a cancellation after pickup must not claim that the text stayed queued");
|
||||
assertEquals(List.of("already sent"), sent());
|
||||
}
|
||||
|
||||
@Test
|
||||
void activeWhileQueuedOrInFlightThenQuietAfterTurnCompletes() {
|
||||
assertTrue(injector.activeTargets().isEmpty());
|
||||
@@ -378,7 +404,7 @@ class InjectorTest {
|
||||
void sendFailureDropsMessageAndFailsItsFuture() {
|
||||
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
|
||||
Injector inj = new Injector(new AgentControl(failing));
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T)).completion();
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
assertTrue(f.isCompletedExceptionally());
|
||||
@@ -387,7 +413,7 @@ class InjectorTest {
|
||||
|
||||
@Test
|
||||
void dropFailsPendingWaiters() {
|
||||
CompletableFuture<Void> f = injector.enqueue(T, "orphan", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> f = injector.enqueue(T, "orphan", TestTurnTokens.inert(T)).completion();
|
||||
injector.drop(T, new HerdrException("worker gone", "pane_not_found", null));
|
||||
assertTrue(f.isCompletedExceptionally(), "queued waiters unblock when the worker vanishes");
|
||||
}
|
||||
@@ -396,8 +422,8 @@ class InjectorTest {
|
||||
void dropPassesTheRealCauseForQueuedAndDeliveredWork() {
|
||||
Captor cap = new Captor();
|
||||
Injector inj = new Injector(new AgentControl(herdr), cap);
|
||||
CompletableFuture<Void> delivered = inj.enqueue(T, "delivered", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> queued = inj.enqueue(T, "queued", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> delivered = inj.enqueue(T, "delivered", TestTurnTokens.inert(T)).completion();
|
||||
CompletableFuture<Void> queued = inj.enqueue(T, "queued", TestTurnTokens.inert(T)).completion();
|
||||
|
||||
inj.onStatus(T, AgentStatus.IDLE); // deliver the first message
|
||||
inj.onStatus(T, AgentStatus.WORKING); // its turn is now in flight; one remains queued
|
||||
@@ -449,7 +475,7 @@ class InjectorTest {
|
||||
Captor cap = new Captor();
|
||||
List<String> forgotten = new ArrayList<>();
|
||||
Injector inj = new Injector(new AgentControl(herdr), cap, _ -> false, forgotten::add);
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "task", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "task", TestTurnTokens.inert(T)).completion();
|
||||
|
||||
for (int i = 0; i < READINESS_SAMPLES; i++) inj.onStatus(T, AgentStatus.IDLE);
|
||||
|
||||
@@ -563,7 +589,7 @@ class InjectorTest {
|
||||
StatusPoller poller = new StatusPoller(new AgentControl(idle), inj, 10);
|
||||
poller.start();
|
||||
try {
|
||||
CompletableFuture<Void> delivered = inj.enqueue(T, "via-poller", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> delivered = inj.enqueue(T, "via-poller", TestTurnTokens.inert(T)).completion();
|
||||
delivered.get(2, TimeUnit.SECONDS); // completes when the poller drives the send
|
||||
} finally {
|
||||
poller.stop();
|
||||
@@ -582,7 +608,7 @@ class InjectorTest {
|
||||
void deliveredFutureCarriesSendFailure() {
|
||||
FakeHerdr failing = new FakeHerdr().agentSendFailsWith("send_failed");
|
||||
Injector inj = new Injector(new AgentControl(failing));
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> f = inj.enqueue(T, "boom", TestTurnTokens.inert(T)).completion();
|
||||
inj.onStatus(T, AgentStatus.IDLE);
|
||||
ExecutionException ex = assertThrows(ExecutionException.class, f::get);
|
||||
assertInstanceOf(HerdrException.class, ex.getCause());
|
||||
|
||||
@@ -41,7 +41,7 @@ class StatusPollerRoutingTest {
|
||||
poller.start();
|
||||
try {
|
||||
CompletableFuture<Void> delivered =
|
||||
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET));
|
||||
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET)).completion();
|
||||
// Must resolve quickly: refining against the WRONG daemon (member) never classifies
|
||||
// out of UNKNOWN, so this would time out under the bug.
|
||||
delivered.get(2, TimeUnit.SECONDS);
|
||||
@@ -65,7 +65,7 @@ class StatusPollerRoutingTest {
|
||||
poller.start();
|
||||
try {
|
||||
CompletableFuture<Void> delivered =
|
||||
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET));
|
||||
injector.enqueue(LEAD_TARGET, "via-poller", TestTurnTokens.inert(LEAD_TARGET)).completion();
|
||||
assertThrows(TimeoutException.class, () -> delivered.get(500, TimeUnit.MILLISECONDS),
|
||||
"a lead target must never be refined from the member daemon's pane content");
|
||||
} finally {
|
||||
|
||||
@@ -299,6 +299,41 @@ class CompositePeerLauncherTest {
|
||||
"two adapter kinds sharing one daemon keep the fallback route");
|
||||
}
|
||||
|
||||
@Test
|
||||
void stopThroughTheSingleDaemonShortcutStillClosesTheTabWhenTheFallbackDelegateUsesPanePlacement() {
|
||||
// fleetd #342: the single-daemon shortcut (spawnedBy empty, herdrDaemonCount()==1) always
|
||||
// routes stop() through delegates.getFirst() — here the claude adapter, configured for
|
||||
// PANE placement (its own profiles never create a dedicated tab). The pane being torn
|
||||
// down here actually belongs to the opencode adapter's TAB placement, sharing the same
|
||||
// herdr daemon — mixing placements is the point: every existing stop-fallback test in this
|
||||
// class configures BOTH adapters as "tab", so usesTabPlacement() was true either way and
|
||||
// the mis-routing never showed.
|
||||
//
|
||||
// Before the fix, HerdrPeerLauncher#stop gated tab resolution on usesTabPlacement() of the
|
||||
// delegate it happened to be called through, so the wrongly-routed (pane-placement) claude
|
||||
// adapter never even looked for a tab to close, and the now-empty tab leaked with nothing
|
||||
// to reap it. The fix (fleetd #342) resolves the pane's real tab unconditionally, so the
|
||||
// decision follows the pane's actual placement rather than the fallback delegate's static
|
||||
// config.
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
FleetConfig.Profile claudePane = new FleetConfig.Profile("claude", "http://gx00.gw:8000",
|
||||
"coder", null, "FLEETD_WORKER_TOKEN", List.of("claude"), "pane", "fleetd-workers",
|
||||
"w #{n}", null, null, null);
|
||||
ClaudeCodeLauncher claude = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of("claude", claudePane), "claude", _ -> null);
|
||||
PeerLauncher composite = new CompositePeerLauncher(List.of(claude, opencodeAdapter(herdr)), "claude");
|
||||
|
||||
// "w9:pW" was never spawned through this composite instance, so spawnedBy has no entry for
|
||||
// it (the same in-memory-cache-miss shape a daemon restart leaves behind) — stop() falls
|
||||
// through to the single-daemon shortcut and hands it to delegates.getFirst() (claude).
|
||||
composite.stop("w9:pW");
|
||||
|
||||
assertTrue(herdr.called("pane.close"), "the pane itself is still closed on every routed path");
|
||||
assertTrue(herdr.called("tab.close"),
|
||||
"the pane's real (sole-occupant) tab must be closed even though the fallback routed "
|
||||
+ "through a delegate configured for pane placement");
|
||||
}
|
||||
|
||||
@Test
|
||||
void listKeepsBothPanesWhenTwoDaemonsShareAPaneId() {
|
||||
// herdr pane ids are per-daemon counters, so two daemons really can both hold w1:p1 on
|
||||
|
||||
@@ -23,6 +23,7 @@ import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
@@ -370,6 +371,161 @@ class HerdrPeerLauncherAllowListWiringTest {
|
||||
"expected the pre-existing 'scrub blanks them' INFO unchanged, got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #341: {@code unprotectedGapLogged} guarded TWO WARN branches that name DIFFERENT env
|
||||
* var names — the allow-list branch ({@code keptByDerivedList}, below) and the deny-by-default
|
||||
* / non-zsh-fallback branch ({@link HerdrPeerLauncher#warnGapUnprotected}). {@code
|
||||
* memberCredentials} is a live, re-read-per-spawn supplier, so the policy can change between
|
||||
* two spawns on the same launcher instance — a config reload needs no restart. Spawn 1 runs
|
||||
* under {@code deny-by-default} with a gap of {@code SPAWN_ONE_UNCOVERED_TOKEN}, which trips
|
||||
* the (before this fix) SHARED one-shot flag. The policy is then reloaded to {@code
|
||||
* allow-list}; spawn 2's gap is {@code FLEETD_WORKER_TOKEN} instead — the test profile's own
|
||||
* {@code tokenEnv}, which the derived allow-list keeps even though it is on neither {@code
|
||||
* known:} nor {@code allow:}, so it is genuinely unprotected and deserves its own WARN. Before
|
||||
* this fix that WARN never fires, because the shared flag was already {@code true} — the
|
||||
* operator is never told {@code FLEETD_WORKER_TOKEN} reaches every member pane unblocked. Real
|
||||
* path: two real {@link HerdrPeerLauncher#spawn} calls on ONE launcher instance, with mutable
|
||||
* {@code memberCredentials}/host-env suppliers standing in for a live config reload between
|
||||
* spawns.
|
||||
*/
|
||||
@Test
|
||||
void aDifferentUnprotectedGapOnALaterSpawnIsNotSuppressedByAnEarlierSpawnsWarn() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AtomicReference<FleetConfig.MemberCredentials> credsState = new AtomicReference<>(
|
||||
new FleetConfig.MemberCredentials(null, List.of(), List.of(), null)); // deny-by-default
|
||||
AtomicReference<Set<String>> hostEnvState =
|
||||
new AtomicReference<>(Set.of("SPAWN_ONE_UNCOVERED_TOKEN"));
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, credsState::get, "/bin/zsh", hostEnvState::get);
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.WARN);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
// Spawn 1: deny-by-default, gap = {SPAWN_ONE_UNCOVERED_TOKEN} — the effectiveAllowed ==
|
||||
// null branch, via warnGapUnprotected.
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
|
||||
// Live policy reload to allow-list, with a DIFFERENT gap name.
|
||||
credsState.set(new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null));
|
||||
hostEnvState.set(Set.of("FLEETD_WORKER_TOKEN"));
|
||||
// Spawn 2: allow-list, gap = {FLEETD_WORKER_TOKEN} — kept by the derived allow-list
|
||||
// (the profile's own tokenEnv), so it is the effectiveAllowed != null / keptByDerivedList
|
||||
// branch, at the SAME log line HerdrPeerLauncher:1819 guards with the shared flag.
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
List<String> messages = appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
assertTrue(messages.stream().anyMatch(
|
||||
m -> m.contains("UNBLOCKED") && m.contains("SPAWN_ONE_UNCOVERED_TOKEN")),
|
||||
"spawn 1's deny-by-default gap must still warn — got: " + messages);
|
||||
assertTrue(messages.stream().anyMatch(
|
||||
m -> m.contains("UNBLOCKED") && m.contains("FLEETD_WORKER_TOKEN")),
|
||||
"spawn 2's gap names a DIFFERENT env var than spawn 1 (FLEETD_WORKER_TOKEN, not "
|
||||
+ "SPAWN_ONE_UNCOVERED_TOKEN) — it must still be warned about even though a "
|
||||
+ "flag already fired once for spawn 1's unrelated name. Before fleetd #341's "
|
||||
+ "fix this WARN never fires because unprotectedGapLogged was already true. "
|
||||
+ "Got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #341 follow-up: the OTHER half of the guard's contract. The set exists to report every
|
||||
* distinct name, but it must still report each one only ONCE — the noise control is the reason
|
||||
* a guard is here at all, and the ticket named it as invariant 1. Two spawns, same policy, same
|
||||
* gap name: exactly one WARN mentioning it.
|
||||
*
|
||||
* <p>Measured before this test existed: replacing {@code .filter(unprotectedGapNamesWarned::add)}
|
||||
* with a filter that adds and always returns {@code true} — so every name is logged on every
|
||||
* spawn — left all 1358 tests green. The fix was correct and nothing held it there. That is the
|
||||
* "a test on the seam does not prove the caller" shape: the {@code Set} behaves, and nothing
|
||||
* proved this class used it as a guard rather than as a record.
|
||||
*/
|
||||
@Test
|
||||
void theSameUnprotectedNameIsWarnedAboutOnlyOnceAcrossSpawns() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AtomicReference<FleetConfig.MemberCredentials> credsState = new AtomicReference<>(
|
||||
new FleetConfig.MemberCredentials(null, List.of(), List.of(), null)); // deny-by-default
|
||||
AtomicReference<Set<String>> hostEnvState =
|
||||
new AtomicReference<>(Set.of("REPEATED_UNCOVERED_TOKEN"));
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, credsState::get, "/bin/zsh", hostEnvState::get);
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.WARN);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
// Same policy, same gap, second spawn. Nothing new to tell the operator.
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
List<String> messages = appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
long mentioning = messages.stream()
|
||||
.filter(m -> m.contains("REPEATED_UNCOVERED_TOKEN"))
|
||||
.count();
|
||||
assertEquals(1, mentioning,
|
||||
"one unchanged unprotected name across two spawns must produce exactly one WARN — "
|
||||
+ "the set is a guard, not just a record. Got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #341 follow-up: the reverse policy order. The original defect was found going
|
||||
* deny-by-default then allow-list, and a guard that is fixed in one direction is not
|
||||
* necessarily fixed in the other — "ask which states still OPEN the gate". Here spawn 1 runs
|
||||
* under {@code allow-list} (the {@code keptByDerivedList} branch) and spawn 2 under
|
||||
* {@code deny-by-default} ({@link HerdrPeerLauncher#warnGapUnprotected}), with a different name
|
||||
* each time. Both must be reported.
|
||||
*/
|
||||
@Test
|
||||
void anAllowListWarnDoesNotSuppressALaterDenyByDefaultWarnForADifferentName() {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
AtomicReference<FleetConfig.MemberCredentials> credsState = new AtomicReference<>(
|
||||
new FleetConfig.MemberCredentials(
|
||||
FleetConfig.MemberCredentials.POLICY_ALLOW_LIST, List.of(), List.of(), null));
|
||||
AtomicReference<Set<String>> hostEnvState =
|
||||
new AtomicReference<>(Set.of("FLEETD_WORKER_TOKEN"));
|
||||
WiringLauncher launcher = new WiringLauncher(herdr, credsState::get, "/bin/zsh", hostEnvState::get);
|
||||
|
||||
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||
Level original = logger.getLevel();
|
||||
logger.setLevel(Level.WARN);
|
||||
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||
appender.start();
|
||||
logger.addAppender(appender);
|
||||
try {
|
||||
// Spawn 1: allow-list, gap kept by the derived list — the keptByDerivedList WARN.
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
|
||||
// Live reload the OTHER way: back to deny-by-default, with a different name.
|
||||
credsState.set(new FleetConfig.MemberCredentials(null, List.of(), List.of(), null));
|
||||
hostEnvState.set(Set.of("LATER_UNCOVERED_TOKEN"));
|
||||
launcher.spawn(new SpawnRequest("test", null, null, null, null, MemberRole.DEV));
|
||||
} finally {
|
||||
logger.detachAppender(appender);
|
||||
logger.setLevel(original);
|
||||
}
|
||||
|
||||
List<String> messages = appender.list.stream().map(ILoggingEvent::getFormattedMessage).toList();
|
||||
assertTrue(messages.stream().anyMatch(
|
||||
m -> m.contains("UNBLOCKED") && m.contains("FLEETD_WORKER_TOKEN")),
|
||||
"spawn 1's allow-list gap must warn — got: " + messages);
|
||||
assertTrue(messages.stream().anyMatch(
|
||||
m -> m.contains("UNBLOCKED") && m.contains("LATER_UNCOVERED_TOKEN")),
|
||||
"spawn 2's deny-by-default gap names a different variable and must still be warned "
|
||||
+ "about, even though an allow-list WARN already fired. Got: " + messages);
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #185 stage 2: with {@code memberHerdrSocket:} configured, member panes run under a
|
||||
* different OS user — {@link HerdrPeerLauncher#hostEnvNames} describes fleetd's own process, not
|
||||
|
||||
@@ -258,7 +258,7 @@ class MessageServiceTest {
|
||||
injector.onStatus(T, AgentStatus.IDLE); // first delivery
|
||||
injector.onStatus(T, AgentStatus.WORKING); // first turn in flight
|
||||
|
||||
CompletableFuture<Void> queued = injector.enqueue(T, "second task", TestTurnTokens.inert(T));
|
||||
CompletableFuture<Void> queued = injector.enqueue(T, "second task", TestTurnTokens.inert(T)).completion();
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(T);
|
||||
injector.drop(T, new HerdrException("agent target sol not found", "agent_not_found", null));
|
||||
|
||||
@@ -377,6 +377,76 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #334. {@code ask()}'s {@code TimeoutException} catch used to forget this task's
|
||||
* {@code turnId} mapping ({@code clearAsyncQuestion(turnId, true)}) BEFORE closing the ask
|
||||
* ({@code rendezvous.closeAsk}, in the shared {@code finally}). A primary's {@code answer()}
|
||||
* call racing that exact window found the ask still "answerable" ({@code
|
||||
* rendezvous.askSession(turnId)} still non-null) while the {@code Task} was already forgotten,
|
||||
* so its {@code task != null} guard skipped the completion and the async ticket sat at
|
||||
* {@code PENDING} forever even though {@code answer()} itself reported a result. The fix
|
||||
* (closing the ask first) makes this window impossible: a racing {@code answer()} call either
|
||||
* still finds the ask open (and the {@code Task} mapping guaranteed intact) or finds it already
|
||||
* closed (and bails {@code STALE_TURN} before ever touching the {@code Task}). This test pins
|
||||
* the exact window with {@code askTimeoutRaceHookForTest} and proves both invariants the ticket
|
||||
* named: (1) a late/racing {@code answer()} sees the ask as already lapsed ({@code STALE_TURN}),
|
||||
* never made answerable again, and (2) the async ticket still resolves {@code DONE} once the
|
||||
* worker's real {@code fleet_reply} lands — it is never stranded {@code PENDING}.
|
||||
*/
|
||||
@Test
|
||||
void aLateAnswerDuringAskTimeoutTeardownStillCompletesTheAsyncTicket() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 200));
|
||||
|
||||
// Wait for the question to actually surface (poll sees ASKING) before racing the timeout.
|
||||
MessageService.TaskView asking = null;
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while ((asking == null || asking.phase() != MessageService.Phase.ASKING)
|
||||
&& System.currentTimeMillis() < deadline) {
|
||||
asking = messages.poll(ticket);
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertNotNull(asking, "the question must surface before the ask times out");
|
||||
String turnId = asking.turnId();
|
||||
assertNotNull(turnId, "an ASKING view carries the turnId to answer on");
|
||||
|
||||
CompletableFuture<MessageService.Reply> lateAnswer = new CompletableFuture<>();
|
||||
messages.setAskTimeoutRaceHookForTest(() ->
|
||||
lateAnswer.complete(messages.answer(turnId, "too late", 500)));
|
||||
try {
|
||||
MessageService.AskResult a = ask.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.AskOutcome.TIMED_OUT, a.outcome());
|
||||
|
||||
MessageService.Reply late = lateAnswer.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.STALE_TURN, late.outcome(),
|
||||
"a late answer racing the timeout teardown must see the ask as already lapsed");
|
||||
|
||||
// The worker resumes on its own (per the ask() contract) and eventually sends its real
|
||||
// fleet_reply; the async ticket must still resolve with it, not strand at PENDING.
|
||||
assertTrue(messages.reply(T, "real result"), "the worker's real reply must still be accepted");
|
||||
} finally {
|
||||
messages.setAskTimeoutRaceHookForTest(null);
|
||||
}
|
||||
|
||||
MessageService.TaskView done = null;
|
||||
deadline = System.currentTimeMillis() + 2000;
|
||||
while ((done == null || done.phase() == MessageService.Phase.PENDING)
|
||||
&& System.currentTimeMillis() < deadline) {
|
||||
done = messages.poll(ticket);
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertNotNull(done);
|
||||
assertEquals(MessageService.Phase.DONE, done.phase(), "the async ticket must not be stranded PENDING");
|
||||
assertEquals("real result", done.reply());
|
||||
}
|
||||
|
||||
@Test
|
||||
void answeringAnUnknownTurnIsStale() {
|
||||
MessageService.Reply r = messages.answer(T + "#999", "too late", 500);
|
||||
@@ -410,6 +480,29 @@ class MessageServiceTest {
|
||||
"a delivered send whose worker never replies times out as still working");
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #345. This forces the injector to pick up the exact pending delivery after {@code send}
|
||||
* first observes its completion as incomplete, but before {@code cancel} takes the target monitor.
|
||||
* The timeout must use {@link Injector.Cancellation#DELIVERED} from {@code cancel} and report
|
||||
* {@link MessageService.Outcome#TIMED_OUT_WORKING}, because the text landed.
|
||||
*
|
||||
* <p>What this does not prove: that this precise interleaving happens by itself under production
|
||||
* timing. The test forces it through a test-only hook; it proves the timeout caller handles the
|
||||
* injector result when the interleaving occurs.
|
||||
*/
|
||||
@Test
|
||||
void sendTimeoutUsesCancellationDeliveredWhenPickupWinsTheRace() {
|
||||
messages.setTimeoutCancellationRaceHookForTest(() -> injector.onStatus(T, AgentStatus.IDLE));
|
||||
try {
|
||||
MessageService.Reply reply = messages.send(T, "race delivery", 50);
|
||||
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_WORKING, reply.outcome(),
|
||||
"cancel reporting DELIVERED means the worker received the timed-out message");
|
||||
} finally {
|
||||
messages.setTimeoutCancellationRaceHookForTest(null);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void answerTimesOutWhenTheResumedWorkerNeverReplies() throws Exception {
|
||||
CompletableFuture<MessageService.Reply> send = sendAsync();
|
||||
@@ -785,6 +878,49 @@ class MessageServiceTest {
|
||||
assertFailedTicket(third, "agent target term_a not found");
|
||||
}
|
||||
|
||||
// --- fleetd #335 (site 1): a per-task cleanup failure inside the abandon() loop must not -----
|
||||
// strand the tasks that come after it. abandon()'s own comment on the loop documents the one
|
||||
// real production call that can throw there (inbox.publish, in the stranded-reply put-back
|
||||
// branch, reached when a concurrent reply() or a second abandon() races this one) — but
|
||||
// reaching that branch requires the exact combination the #137 follow-up above already found
|
||||
// unreachable through the public API. abandonCleanupHookForTest reproduces the resulting SHAPE
|
||||
// (one task's cleanup throws) directly instead, the same technique this file already uses for
|
||||
// fleetd #324/#329's own hard-to-time races.
|
||||
@Test
|
||||
void aPerTaskCleanupFailureDoesNotStrandTheRemainingMatchingTasks() throws Exception {
|
||||
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
|
||||
try {
|
||||
String first = messages.sendAsync(T, "first task");
|
||||
awaitWaiting(); // first task owns the target lock and rendezvous waiter
|
||||
String second = messages.sendAsync(T, "second task"); // parked on the same lock
|
||||
String third = messages.sendAsync(T, "third task"); // parked too — the whole sweep must survive
|
||||
|
||||
java.util.concurrent.atomic.AtomicInteger calls = new java.util.concurrent.atomic.AtomicInteger();
|
||||
messages.setAbandonCleanupHookForTest(() -> {
|
||||
if (calls.getAndIncrement() == 0) {
|
||||
throw new RuntimeException("PROBE-335-SITE1");
|
||||
}
|
||||
});
|
||||
|
||||
assertTrue(messages.abandon(T, "agent target term_a not found"));
|
||||
|
||||
// Every task in the loop still gets its own outcome — the one whose cleanup threw
|
||||
// included — even though the loop had no way to know in advance which one that would be.
|
||||
assertFailedTicket(first, "agent target term_a not found");
|
||||
assertFailedTicket(second, "agent target term_a not found");
|
||||
assertFailedTicket(third, "agent target term_a not found");
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e ->
|
||||
e.getLevel() == Level.ERROR
|
||||
&& e.getThrowableProxy() != null
|
||||
&& "PROBE-335-SITE1".equals(e.getThrowableProxy().getMessage())),
|
||||
"a per-task cleanup failure must still reach the log, not vanish silently");
|
||||
} finally {
|
||||
messages.setAbandonCleanupHookForTest(null);
|
||||
detachMessageServiceLog(appender);
|
||||
}
|
||||
}
|
||||
|
||||
// --- #137 follow-up: abandon() must not guess when more than one task is open ---------------
|
||||
//
|
||||
// A test combining a genuine stranded reply (hasStrandedReply(T)==true) with two simultaneously
|
||||
@@ -1472,6 +1608,42 @@ class MessageServiceTest {
|
||||
}
|
||||
}
|
||||
|
||||
// --- fleetd #335 (site 2): task.future.whenComplete's own returned stage is discarded, so an --
|
||||
// uncaught throw from ReplyPushLoop.onTicketTerminal used to vanish with no log line and no
|
||||
// metric. The real production trigger is the daemon's own shutdown sequence (Fleetd's shutdown
|
||||
// hook): messages.close() only stops the async executor from taking NEW work — it does not
|
||||
// cancel a send already in flight — while pushLoop.close() shuts its scheduler down immediately
|
||||
// right after, so a ticket that completes in that narrow window has onTicketTerminal's own
|
||||
// scheduler.schedule(...) throw a real RejectedExecutionException. Reproduced here by shutting
|
||||
// the very same scheduler down before the ticket resolves — no test-only hook needed, this
|
||||
// reachable path throws for real.
|
||||
@Test
|
||||
void aTicketTerminalPushFailureDoesNotVanishSilently() throws Exception {
|
||||
ListAppender<ILoggingEvent> appender = attachMessageServiceLog();
|
||||
try (var wiring = wireWithPushLoop(1, 50)) {
|
||||
String ticket = wiring.service().sendAsync(T, "long task");
|
||||
awaitWaiting();
|
||||
wiring.scheduler().shutdownNow(); // simulate pushLoop.close() racing an in-flight send
|
||||
injectDelivery();
|
||||
assertTrue(rendezvous.resolve(T, "async result"));
|
||||
|
||||
// The ticket's own outcome must be unaffected by the swallowed exception — finishAsyncTask
|
||||
// completes task.future before whenComplete's action (and thus onTicketTerminal) ever runs.
|
||||
MessageService.TaskView done = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.DONE);
|
||||
assertEquals("async result", done.reply());
|
||||
|
||||
assertTrue(appender.list.stream().anyMatch(e ->
|
||||
e.getLevel() == Level.ERROR
|
||||
&& e.getFormattedMessage().contains(ticket)
|
||||
&& e.getThrowableProxy() != null
|
||||
&& "java.util.concurrent.RejectedExecutionException"
|
||||
.equals(e.getThrowableProxy().getClassName())),
|
||||
"onTicketTerminal throwing must still reach the log, not vanish silently");
|
||||
} finally {
|
||||
detachMessageServiceLog(appender);
|
||||
}
|
||||
}
|
||||
|
||||
// --- CB-582: fleet_ask question-open nudges --------------------------------------------------
|
||||
|
||||
@Test
|
||||
@@ -1768,7 +1940,10 @@ class MessageServiceTest {
|
||||
assertEquals(MessageService.Outcome.TIMED_OUT_QUEUED, r.outcome());
|
||||
|
||||
assertTrue(messages.hasQueuedDelivery(T),
|
||||
"a TIMED_OUT_QUEUED send leaves the message still queued in the injector");
|
||||
"a TIMED_OUT_QUEUED send still records the undelivered delivery for fleet health");
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
assertTrue(herdr.calls.stream().noneMatch(c -> c.method().equals("agent.prompt")),
|
||||
"a TIMED_OUT_QUEUED send must be cancelled, not delivered when the worker later goes idle");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -711,13 +711,21 @@ class FleetAppTest {
|
||||
|
||||
@Test
|
||||
void stopWorkerInPanePlacementClosesOnlyThePane() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
// fleetd #342: tab-cleanup resolution is no longer skipped based on a profile's declared
|
||||
// placement — a stop() routed through the wrong delegate (e.g. CompositePeerLauncher's
|
||||
// single-daemon shortcut after a daemon restart) could carry a placement config that says
|
||||
// nothing true about how THIS pane was actually placed. So the pane's real tab is now
|
||||
// always resolved, and the single-occupant check below is what protects a pane-placement
|
||||
// peer's shared tab, exactly as it always protected a tab-placement one. Model that
|
||||
// realistically: the peer's pane was split into an existing tab that already held another
|
||||
// occupant, so the tab must never be closed.
|
||||
FakeHerdr herdr = new FakeHerdr().withWorkerTabPaneCount(2);
|
||||
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"), "pane");
|
||||
|
||||
assertEquals(204, req(port, "DELETE", "/members/w9:pW").statusCode());
|
||||
assertTrue(herdr.called("pane.close"));
|
||||
assertFalse(herdr.called("tab.close"), "pane placement owns no tab to close");
|
||||
assertFalse(herdr.called("pane.get"), "no tab resolution in pane placement");
|
||||
assertTrue(herdr.called("pane.get"), "tab resolution now always runs, regardless of placement");
|
||||
assertFalse(herdr.called("tab.close"), "pane placement's shared tab must never be closed");
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user