Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 6cccd458d4 |
@@ -1,23 +0,0 @@
|
||||
---
|
||||
name: hunter
|
||||
description: Sweep one assigned scope for defects and report ranked findings without changes.
|
||||
---
|
||||
|
||||
<!-- CB-617: The model comes from fleetd.yaml because the launch flag overrides model here on both backends. -->
|
||||
|
||||
You sweep the assigned package or scope for real defects. Read the full assigned scope before you
|
||||
judge it. Report several ranked findings when the evidence supports them. Change nothing: do not
|
||||
edit code, commit, push, or open a pull request.
|
||||
|
||||
You may run the build or tests to check a finding. Read the complete output and report the real
|
||||
result. Do not hide failures with a pipe. State only checks you actually ran. The primary's IDE
|
||||
tools are not yours. A mounted forge tool may use a blocked credential and fail by design.
|
||||
|
||||
Do only the assigned scope. Note anything outside it in one line and do not investigate it further.
|
||||
Use `fleet_ask{question}` only when a decision belongs to the lead, such as an unclear requirement
|
||||
or two defensible fixes. Do not ask about something you can decide by reading more code.
|
||||
|
||||
Your handoff must name the files you read, each ranked finding or `NO FINDINGS`, the checks you ran,
|
||||
and any caveat for review.
|
||||
|
||||
The launcher provides the required bridge reply instructions for every member.
|
||||
@@ -284,8 +284,6 @@ must obey belongs in the charter, not here.
|
||||
it for a multi-finding sweep hands the worker two contradictory output contracts. That has
|
||||
already cost three workers' turns: each wrote a good report to its terminal and ended the turn
|
||||
with no `fleet_reply`, and the scrape returned the tail of the brief instead.
|
||||
Spawn `implementer` with role `dev`, `reviewer` with role `reviewer`, and `hunter` with role
|
||||
`hunter`.
|
||||
- **Primary-side skills** (not delegation playbooks — a worker cannot use them):
|
||||
`port-to-opencode` (make an OpenCode session a participant in this workspace),
|
||||
`fleets-status` (report every fleet that shares one LavinMQ instance),
|
||||
|
||||
@@ -560,7 +560,7 @@ placement: weighted
|
||||
# older keys: `leaders:`, `members:`, `leadScan:` and `defaultProfile:`.
|
||||
#
|
||||
# A member is anything a lead spawns, and every member has two INDEPENDENT attributes:
|
||||
# role — which contract: architect, dev, hunter or reviewer. It picks the launch charter, the role
|
||||
# role — which contract: architect, dev or reviewer. It picks the launch charter, the role
|
||||
# file, the playbook skill and the authz row.
|
||||
# profile — which backend: one of the `profiles:` keys above (model, CLI adapter, cost).
|
||||
# They vary on their own. A reviewer may run on the same profile as the dev whose diff it reads,
|
||||
@@ -572,13 +572,13 @@ placement: weighted
|
||||
#
|
||||
# Each pool lists the profiles that role MAY run on — these are pools, not identities. That is also
|
||||
# what replaced `defaultProfile:`: an unqualified spawn names a role, and that role's pool supplies
|
||||
# the candidates, in definition order. A dev, hunter and reviewer staying anonymous is exactly
|
||||
# compatible with being listed here; the entry key just names the entry.
|
||||
# the candidates, in definition order. A dev and a reviewer staying anonymous is exactly compatible
|
||||
# with being listed here; the entry key just names the entry.
|
||||
fleet:
|
||||
# Optional launch-charter text, keyed only by the singular role wire names: architect, dev,
|
||||
# hunter, reviewer. Changes are HOT and reach the next spawn without a daemon restart. Do not put
|
||||
# secrets here: a later launch step writes this text to a world-readable temp file, and ${ENV}
|
||||
# interpolation is deliberately not supported.
|
||||
# reviewer. Changes are HOT and reach the next spawn without a daemon restart. Do not put secrets
|
||||
# here: a later launch step writes this text to a world-readable temp file, and ${ENV} interpolation
|
||||
# is deliberately not supported.
|
||||
charters:
|
||||
architect: |-
|
||||
You are an architect in this fleet. You refine work before anyone builds it:
|
||||
@@ -589,9 +589,6 @@ fleet:
|
||||
dev: |-
|
||||
You implement the one unit you were given, and nothing else. You test it,
|
||||
commit it, and open your own pull request. You never merge.
|
||||
hunter: |-
|
||||
You sweep the assigned scope for real defects. You may run the build or tests
|
||||
to check a finding. You change nothing, and report several ranked findings.
|
||||
reviewer: |-
|
||||
You review the diff you were given. You report bugs, risks and missing tests.
|
||||
You do not change code.
|
||||
@@ -669,9 +666,6 @@ fleet:
|
||||
developers:
|
||||
gx10:
|
||||
profile: gx10
|
||||
# hunters:
|
||||
# gx10:
|
||||
# profile: gx10 # a hunt may run checks, but never changes code
|
||||
# reviewers:
|
||||
# gx10:
|
||||
# profile: gx10 # the same backend may serve two roles; that is the point
|
||||
|
||||
@@ -24,6 +24,7 @@ import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.inject.LoopWatchdog;
|
||||
import dev.ltms.fleet.inject.StatusPoller;
|
||||
import dev.ltms.fleet.inject.TurnListener;
|
||||
import dev.ltms.fleet.inject.TurnRegistrar;
|
||||
import dev.ltms.fleet.inject.MemberPresence;
|
||||
import dev.ltms.fleet.auth.MemberRegistry;
|
||||
import dev.ltms.fleet.auth.CallerResolver;
|
||||
@@ -80,7 +81,9 @@ import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
import java.util.function.BiConsumer;
|
||||
import java.util.function.BooleanSupplier;
|
||||
import java.util.function.Consumer;
|
||||
import java.util.function.Function;
|
||||
import java.util.function.LongSupplier;
|
||||
import java.util.function.Predicate;
|
||||
@@ -501,21 +504,26 @@ public final class Fleetd {
|
||||
// fleetd #556: registration is wired directly to `completion`, not folded into the
|
||||
// `turnListener` fan-out above — so it survives `sessions.onDelivered` (or any future
|
||||
// listener) throwing, regardless of call order. See TurnRegistrar's javadoc.
|
||||
// fleetd #589 Group 3 (:505): extracted to turnRegistrar(...) — see FleetdTurnRegistrarWiringTest.
|
||||
Injector injector = new Injector(router, turnListener, deliverable,
|
||||
presence::forget, completion::register);
|
||||
presence::forget, turnRegistrar(completion));
|
||||
StatusPoller poller = new StatusPoller(router, injector, Injector.POLL_INTERVAL_MILLIS);
|
||||
poller.start();
|
||||
|
||||
// CB-307: reply inbox. A broker: block selects the AMQP-backed durable adapter; absent (or
|
||||
// unusable), fleetd stays soft-state on the in-memory inbox. The AMQP inbox owns a broker
|
||||
// connection, so keep the reference to close it in the ordered shutdown hook.
|
||||
final ReplyInbox replyInbox = selectReplyInbox(cfg.broker(), System.getenv(), AmqpReplyInbox::open);
|
||||
// fleetd #589 Group 3 (:512): opener extracted to replyInboxOpener() — see
|
||||
// FleetdReplyInboxOpenerWiringTest.
|
||||
final ReplyInbox replyInbox = selectReplyInbox(cfg.broker(), System.getenv(), replyInboxOpener());
|
||||
// CB-637: this daemon's lead-to-lead mailbox on the SHARED coordination vhost — a separate
|
||||
// broker from the reply inbox by design (see FleetConfig.Coordinator). Absent a coordinator:
|
||||
// block this is null and every lead path below is simply not wired, which is exactly the
|
||||
// behaviour before this ticket. It owns a broker connection, so keep the reference for the
|
||||
// ordered shutdown hook.
|
||||
final LeadMailbox leadMailbox = openLeadMailbox(cfg.coordinator(), System.getenv(), LeadMailbox::open);
|
||||
// fleetd #589 Group 3 (:518): opener extracted to leadMailboxOpener() — see
|
||||
// FleetdLeadMailboxOpenerWiringTest.
|
||||
final LeadMailbox leadMailbox = openLeadMailbox(cfg.coordinator(), System.getenv(), leadMailboxOpener());
|
||||
// CB-307: learn the primary's terminal from orchestration tool calls (or pin from config).
|
||||
// The pin also feeds CallerResolver below: a primary running inside a herdr pane would
|
||||
// otherwise resolve as a worker and be refused every orchestration tool.
|
||||
@@ -585,10 +593,12 @@ public final class Fleetd {
|
||||
// fleetd #386: System::nanoTime freezes across a macOS sleep, so the stall check also
|
||||
// gets a wall-clock source to detect and correct for that freeze. Every other decision
|
||||
// in FleetHealthMonitor stays on the monotonic clock, unchanged.
|
||||
// fleetd #589 Group 3 (:591): failTarget extracted to healthFailTarget(...) — see
|
||||
// FleetdHealthFailTargetWiringTest.
|
||||
healthMonitor = new FleetHealthMonitor(agents, sessions::roster, messages, healthScheduler,
|
||||
System::nanoTime, () -> TimeUnit.MILLISECONDS.toNanos(System.currentTimeMillis()),
|
||||
cfg.health().intervalOrDefault(),
|
||||
cfg.health().workingSuspectAfterOrDefault(), messages::abandon);
|
||||
cfg.health().workingSuspectAfterOrDefault(), healthFailTarget(messages));
|
||||
String coverage = FleetHealthMonitor.coverage(true,
|
||||
cfg.health().notifications() != null && cfg.health().notifications().configured());
|
||||
if ("detection-only".equals(coverage)) {
|
||||
@@ -608,27 +618,11 @@ public final class Fleetd {
|
||||
// CB-516: releasing a worker must fail whatever send was waiting on it. Without this a
|
||||
// torn-down delegation kept reporting PENDING until the 30-minute async timeout, and never
|
||||
// reached /metrics — the delegation was unresolvable and nothing said so.
|
||||
sessions.onRelease(detail -> {
|
||||
// CB-578 stage C, acceptance criterion 10: a failed ticket's detail should tell a lead
|
||||
// where to re-dispatch onto the same tree, not just that the worker vanished.
|
||||
String reason = "the worker session was released before it replied";
|
||||
if (detail.worktreePath() != null) {
|
||||
reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch()
|
||||
+ " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none");
|
||||
}
|
||||
// CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the
|
||||
// member's conversation instead of only re-dispatching a fresh one onto the same files.
|
||||
if (detail.agentSessionId() != null) {
|
||||
reason += " agentSessionId=" + detail.agentSessionId();
|
||||
}
|
||||
// fleetd #275: this is an explicit teardown (fleet_stop, or the idle reaper) — the
|
||||
// worker's pane is being stopped right now, so an open fleet_ask has no turn left to
|
||||
// resume into. Sweep it too, unlike FleetHealthMonitor's health-classification call
|
||||
// (see MessageService.abandon's javadoc for why those two must differ).
|
||||
messages.abandon(detail.terminalId(), reason, true);
|
||||
replyInbox.release(detail.terminalId());
|
||||
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
|
||||
});
|
||||
// fleetd #589 Group 3 (:611-631): the whole cleanup lambda extracted to releaseCleanup(...)
|
||||
// — see FleetdReleaseCleanupWiringTest, and MessageService.abandon's javadoc for the
|
||||
// documented incident (a torn-down worker's rendezvous waiter left open) this lambda exists
|
||||
// to prevent.
|
||||
sessions.onRelease(releaseCleanup(messages, replyInbox, primaryRegistry));
|
||||
|
||||
// MCP server face (CB-105): fleet_send/fleet_reply/fleet_status, mounted at /mcp.
|
||||
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
|
||||
@@ -1066,6 +1060,104 @@ public final class Fleetd {
|
||||
() -> reaper == null ? LoopWatchdog.State.STOPPED : reaper.health());
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :505}): package-private factory for the {@link Injector}'s
|
||||
* {@link TurnRegistrar}, extracted out of {@code main} for the same reason {@link
|
||||
* #loopHealthSource} was — before this ticket {@code completion::register} was an inline
|
||||
* argument to {@code new Injector(...)}, so nothing could pin it directly. Replacing it with
|
||||
* {@link TurnRegistrar#NOOP} compiles clean and leaves every existing test green: {@code
|
||||
* onDelivered}'s own {@code captureBaseline} does the identical {@code inFlight} check-and-put a
|
||||
* moment later on the ordinary path, so the two are indistinguishable once a turn's delivery
|
||||
* finishes normally. The gap {@link TurnRegistrar}'s own javadoc (fleetd #556) exists to close is
|
||||
* a {@code turnListener} callback throwing between the two — {@link
|
||||
* FleetdTurnRegistrarWiringTest} pins that {@code register} itself (not just {@code
|
||||
* captureBaseline}) makes a delivered turn's waiter resolvable.
|
||||
*/
|
||||
static TurnRegistrar turnRegistrar(CompletionResolver completion) {
|
||||
return completion::register;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :591}): package-private factory for {@link
|
||||
* FleetHealthMonitor}'s {@code failTarget} callback, extracted out of {@code main} for the same
|
||||
* reason {@link #loopHealthSource} was. Before this ticket {@code messages::abandon} was an
|
||||
* inline argument to {@code new FleetHealthMonitor(...)}; replacing it with a no-op {@code
|
||||
* BiConsumer} compiles clean and leaves every existing test green, and in production it means a
|
||||
* member found {@code GONE}/{@code NEVER_READY} never fails the ticket waiting on it — the
|
||||
* caller reports {@code PENDING} for the full 30-minute async timeout instead of the immediate,
|
||||
* accurate failure CB-580 exists to give it. {@link FleetdHealthFailTargetWiringTest} pins that
|
||||
* the returned callback actually reaches the real {@link MessageService#abandon}.
|
||||
*/
|
||||
static BiConsumer<String, String> healthFailTarget(MessageService messages) {
|
||||
return messages::abandon;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :611-631}): package-private factory for the whole {@link
|
||||
* SessionManager#onRelease} cleanup callback, extracted out of {@code main} for the same reason
|
||||
* {@link #loopHealthSource} was. Before this ticket this was an inline lambda built directly
|
||||
* inside {@code main}; replacing its body with a no-op {@code detail -> { }} compiles clean and
|
||||
* leaves every existing test green, and in production it is the exact incident {@link
|
||||
* MessageService#abandon}'s own javadoc documents: a torn-down worker's rendezvous waiter is
|
||||
* left open, so a blocking {@code fleet_send} keeps blocking and an async one reports {@code
|
||||
* PENDING} for a hardcoded thirty minutes on every {@code fleet_stop} and every idle-reap.
|
||||
*
|
||||
* <p>{@link FleetdReleaseCleanupWiringTest} pins all three collaborator calls this lambda makes
|
||||
* — {@code messages.abandon}, {@code replyInbox.release}, and {@code
|
||||
* primaryRegistry.forgetDelegation} — each already tested on its own ({@code MessageServiceTest},
|
||||
* {@code PrimaryRegistryTest}), but never before proven to actually be reached from here.
|
||||
*/
|
||||
static Consumer<SessionManager.ReleaseDetail> releaseCleanup(MessageService messages, ReplyInbox replyInbox,
|
||||
PrimaryRegistry primaryRegistry) {
|
||||
return detail -> {
|
||||
// CB-578 stage C, acceptance criterion 10: a failed ticket's detail should tell a lead
|
||||
// where to re-dispatch onto the same tree, not just that the worker vanished.
|
||||
String reason = "the worker session was released before it replied";
|
||||
if (detail.worktreePath() != null) {
|
||||
reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch()
|
||||
+ " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none");
|
||||
}
|
||||
// CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the
|
||||
// member's conversation instead of only re-dispatching a fresh one onto the same files.
|
||||
if (detail.agentSessionId() != null) {
|
||||
reason += " agentSessionId=" + detail.agentSessionId();
|
||||
}
|
||||
// fleetd #275: this is an explicit teardown (fleet_stop, or the idle reaper) — the
|
||||
// worker's pane is being stopped right now, so an open fleet_ask has no turn left to
|
||||
// resume into. Sweep it too, unlike FleetHealthMonitor's health-classification call
|
||||
// (see MessageService.abandon's javadoc for why those two must differ).
|
||||
messages.abandon(detail.terminalId(), reason, true);
|
||||
replyInbox.release(detail.terminalId());
|
||||
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :512}): package-private factory for {@code selectReplyInbox}'s
|
||||
* production {@link AmqpOpener}, extracted out of {@code main} for the same reason {@link
|
||||
* #loopHealthSource} was. Before this ticket {@code AmqpReplyInbox::open} was an inline argument
|
||||
* to the {@code selectReplyInbox(...)} call; replacing it with {@code (uri, prefetch) -> new
|
||||
* InMemoryReplyInbox()} compiles clean and leaves every existing test green — {@link
|
||||
* FleetdReplyInboxSelectionTest} drives {@code selectReplyInbox} with its own injected opener and
|
||||
* never sees what {@code main} actually passes. {@link FleetdReplyInboxOpenerWiringTest} pins
|
||||
* that this returns the real opener by pointing it at a guaranteed-closed local port and
|
||||
* asserting the real network attempt throws — the inert stub never attempts a connection at all.
|
||||
*/
|
||||
static AmqpOpener replyInboxOpener() {
|
||||
return AmqpReplyInbox::open;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3 (site {@code :518}): package-private factory for {@code
|
||||
* openLeadMailbox}'s production {@link LeadMailboxOpener}, extracted out of {@code main} for the
|
||||
* same reason {@link #replyInboxOpener} was — same gap, same fix, the lead-coordination mailbox
|
||||
* instead of the reply inbox. {@link FleetdLeadMailboxOpenerWiringTest} pins that this returns
|
||||
* the real opener the same way.
|
||||
*/
|
||||
static LeadMailboxOpener leadMailboxOpener() {
|
||||
return LeadMailbox::open;
|
||||
}
|
||||
|
||||
/**
|
||||
* fleetd #248: package-private factory for the member worktree/branch lookup {@link
|
||||
* CompletionResolver} uses to name a fallback report's worktree and branch (fleetd#241).
|
||||
|
||||
@@ -221,7 +221,7 @@ public final class CallerResolver {
|
||||
// The config/live binding names this pane as an architect slot's own. Same
|
||||
// unforgeable pane mapping; the live binding, never a request argument, decides.
|
||||
// Check the slot role too: this defence in depth prevents a bad lifecycle bind from
|
||||
// escalating a dev, hunter, or reviewer into an architect. Checked before the worker fallback.
|
||||
// escalating a dev or reviewer into an architect. Checked before the worker fallback.
|
||||
return Principal.architect(memberSlotNames.apply(slot), c.terminal(), c.pid());
|
||||
}
|
||||
return Principal.worker(c.terminal(), c.pid()); // unforgeable; never token-gated
|
||||
|
||||
@@ -42,7 +42,7 @@ public interface MemberLifecycle {
|
||||
* Try to bind a newly spawned {@code terminal} into the role it was granted.
|
||||
*
|
||||
* @return the role this session actually holds: {@code role} unchanged for a role with no
|
||||
* live slot-binding semantics (dev, hunter, reviewer), or when the bind succeeded; a fallback
|
||||
* live slot-binding semantics (dev, reviewer), or when the bind succeeded; a fallback
|
||||
* role — never {@code role} — when a slot-bound role (architect) could not be bound.
|
||||
* Callers must record THIS value on the session, never the requested {@code role}, so
|
||||
* a later roster read never reports a role the session does not hold (CB-619). In
|
||||
|
||||
@@ -20,8 +20,7 @@ import java.util.function.Supplier;
|
||||
*
|
||||
* <p>Two halves, split by who owns each:
|
||||
* <ul>
|
||||
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/{@code hunters}/
|
||||
* {@code reviewers}
|
||||
* <li><b>slots</b> — read from {@code fleet.architects}/{@code developers}/{@code reviewers}
|
||||
* (see {@link #slots()}), each carrying the {@code profile} reference the spawn lifecycle
|
||||
* reads when it stands the slot up. <strong>Live, since fleetd #424</strong>: {@link #live}
|
||||
* re-reads {@code fleet:} on every call, through a supplier the same shape as
|
||||
@@ -323,7 +322,7 @@ public final class MemberRegistry implements MemberLifecycle {
|
||||
* CB-619 / fleetd #123: refuse an architect acquire before anything spawns when no configured
|
||||
* slot carries {@code profile} — the config-gap case from the original defect report (a spawn
|
||||
* asked for {@code role=architect, profile=sonnet}, and {@code fleet.architects} carried only
|
||||
* {@code opus} and {@code sol}). A dev/hunter/reviewer acquire is always a no-op: those pools are
|
||||
* {@code opus} and {@code sol}). A dev/reviewer acquire is always a no-op: those pools are
|
||||
* placement candidates only (see {@code CompositePeerLauncher}), never a live identity binding,
|
||||
* so there is nothing here to refuse — an explicit profile outside the pool for those roles is a
|
||||
* documented operator override, not a defect.
|
||||
|
||||
@@ -31,7 +31,7 @@ import java.util.function.Supplier;
|
||||
* {@code placement:}, and an existing profile's {@code weight} / {@code maxLoad}. Both are
|
||||
* read through a supplier on {@code CompositePeerLauncher}, which is what makes them hot —
|
||||
* not the fact that they are config. Most of {@code fleet:} — every role pool
|
||||
* ({@code architects}/{@code developers}/{@code hunters}/{@code reviewers}), {@code charters}, and
|
||||
* ({@code architects}/{@code developers}/{@code reviewers}), {@code charters}, and
|
||||
* {@code tabLabel} — is read the same live way, through the same supplier
|
||||
* ({@code () -> config.get().fleet()}). {@code architects} in particular is hot for
|
||||
* <strong>two independent consumers</strong> (fleetd #424): {@code CompositePeerLauncher}
|
||||
@@ -605,7 +605,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
+ "opened once and needs a restart; the broker URI env-var name kept out of a "
|
||||
+ "member's environment is read live on every spawn and already applied");
|
||||
}
|
||||
// fleetd #333: unlike health/coordinator above, most of `fleet:` (developers, hunters, reviewers,
|
||||
// fleetd #333: unlike health/coordinator above, most of `fleet:` (developers, reviewers,
|
||||
// charters, tabLabel) is genuinely hot — ConfigRefTest.aHotChangeIsAppliedAndRead-
|
||||
// ThroughGet and aCharterChangeIsHotAndReachesTheLiveConfig prove it reaches the live config
|
||||
// with no restart note. `architects` is hot too, and — since fleetd #424 — hot for BOTH of
|
||||
@@ -630,7 +630,7 @@ public final class ConfigRef implements Supplier<FleetConfig> {
|
||||
+ "identity map and to auto-launch leads, and neither is rebuilt on reload, so a "
|
||||
+ "lead added, removed, or given a new tab: label needs a restart — until then it "
|
||||
+ "stays unrecognised, and a caller from its new tab resolves as a worker, not a "
|
||||
+ "lead; the rest of fleet: (developers, hunters, reviewers, charters, tabLabel) is read "
|
||||
+ "lead; the rest of fleet: (developers, reviewers, charters, tabLabel) is read "
|
||||
+ "live through the supplier on CompositePeerLauncher, and architects is read "
|
||||
+ "live through that same supplier for placement AND through a separate supplier "
|
||||
+ "on MemberRegistry for spawn-time identity — both already applied");
|
||||
|
||||
@@ -62,7 +62,7 @@ import java.util.regex.PatternSyntaxException;
|
||||
* @param fleet who the daemon may run and under which role (CB-557). One block replacing
|
||||
* the former {@code leaders:}, {@code members:}, {@code leadScan:} and
|
||||
* {@code defaultProfile:}. Role is the containing key — {@code leaders},
|
||||
* {@code architects}, {@code developers}, {@code hunters}, {@code reviewers} — and each entry
|
||||
* {@code architects}, {@code developers}, {@code reviewers} — and each entry
|
||||
* names the {@code profiles:} backend it runs on. See {@link Fleet}
|
||||
* @param leadHeartbeat opt-in idle-lead heartbeat (CB-551); {@code null} ⇒ off, and an upgraded
|
||||
* daemon never nudges an idle lead on its own initiative
|
||||
@@ -1200,7 +1200,6 @@ public record FleetConfig(
|
||||
* @param leaders panes that orchestrate rather than are orchestrated, keyed by lead name
|
||||
* @param architects profiles the {@code architect} role may run on
|
||||
* @param developers profiles the {@code dev} role may run on
|
||||
* @param hunters profiles the {@code hunter} role may run on
|
||||
* @param reviewers profiles the {@code reviewer} role may run on
|
||||
* @param charters optional launch-charter text keyed by singular role wire name
|
||||
* @param tabLabel template for a member tab's label; {@code {role}}, {@code {profile}},
|
||||
@@ -1211,7 +1210,6 @@ public record FleetConfig(
|
||||
public record Fleet(Map<String, Leader> leaders,
|
||||
Map<String, Slot> architects,
|
||||
Map<String, Slot> developers,
|
||||
Map<String, Slot> hunters,
|
||||
Map<String, Slot> reviewers,
|
||||
Map<String, String> charters,
|
||||
String tabLabel) {
|
||||
@@ -1228,7 +1226,6 @@ public record FleetConfig(
|
||||
leaders = unmodifiableOrEmpty(leaders);
|
||||
architects = unmodifiableOrEmpty(architects);
|
||||
developers = unmodifiableOrEmpty(developers);
|
||||
hunters = unmodifiableOrEmpty(hunters);
|
||||
reviewers = unmodifiableOrEmpty(reviewers);
|
||||
charters = unmodifiableOrEmpty(charters);
|
||||
tabLabel = (tabLabel == null || tabLabel.isBlank()) ? DEFAULT_TAB_LABEL : tabLabel;
|
||||
@@ -1243,15 +1240,9 @@ public record FleetConfig(
|
||||
* constructor: the launcher reads {@code fleet.charters()} from the live config. Jackson
|
||||
* binds the canonical constructor, so this one cannot swallow an operator's YAML.
|
||||
*/
|
||||
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
|
||||
Map<String, Slot> developers, Map<String, Slot> reviewers,
|
||||
Map<String, String> charters, String tabLabel) {
|
||||
this(leaders, architects, developers, null, reviewers, charters, tabLabel);
|
||||
}
|
||||
|
||||
public Fleet(Map<String, Leader> leaders, Map<String, Slot> architects,
|
||||
Map<String, Slot> developers, Map<String, Slot> reviewers, String tabLabel) {
|
||||
this(leaders, architects, developers, null, reviewers, null, tabLabel);
|
||||
this(leaders, architects, developers, reviewers, null, tabLabel);
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -1273,7 +1264,6 @@ public record FleetConfig(
|
||||
return switch (role) {
|
||||
case ARCHITECT -> architects;
|
||||
case DEV -> developers;
|
||||
case HUNTER -> hunters;
|
||||
case REVIEWER -> reviewers;
|
||||
};
|
||||
}
|
||||
@@ -1882,7 +1872,7 @@ public record FleetConfig(
|
||||
|
||||
/** The {@code fleet:} child blocks whose direct children are slot names. */
|
||||
private static final Set<String> FLEET_POOL_KEYS =
|
||||
Set.of("leaders", "architects", "developers", "hunters", "reviewers");
|
||||
Set.of("leaders", "architects", "developers", "reviewers");
|
||||
|
||||
/**
|
||||
* Reject a {@code fleet:} role pool whose slot names repeat (CB-548, re-homed by CB-557).
|
||||
@@ -1892,7 +1882,7 @@ public record FleetConfig(
|
||||
* daemon would never know. Jackson's YAML parser does not fail on duplicate mapping keys by
|
||||
* default, so duplicates are caught here, at parse time, before the map is built.
|
||||
*
|
||||
* <p>Only the five pools <em>directly under the top-level {@code fleet:}</em> are considered,
|
||||
* <p>Only the four pools <em>directly under the top-level {@code fleet:}</em> are considered,
|
||||
* and only their direct child keys (the slot names). A nested field elsewhere, even one also
|
||||
* named {@code developers:}, is ignored, so parsing of the rest of the config is unaffected.
|
||||
*
|
||||
@@ -2060,7 +2050,7 @@ public record FleetConfig(
|
||||
+ " and that role's pool supplies the candidate profiles",
|
||||
"architects", "'fleet.architects'",
|
||||
"members", "a role pool under 'fleet:' — 'fleet.architects', 'fleet.developers' or"
|
||||
+ " 'fleet.hunters' or 'fleet.reviewers'; the role is the containing key, not a 'role:' field",
|
||||
+ " 'fleet.reviewers'; the role is the containing key, not a 'role:' field",
|
||||
"leaders", "'fleet.leaders'",
|
||||
"leadScan", "'fleet.leaders.<name>.tabPrefix' and '.scanIntervalSeconds' — lead"
|
||||
+ " discovery is now configured on the lead it discovers");
|
||||
|
||||
@@ -2134,9 +2134,8 @@ public final class FleetMcp {
|
||||
return tool(FleetTool.SPAWN.wireName(),
|
||||
"Spawn a new off-subscription member session. A member has two independent attributes: "
|
||||
+ "role (what it is for) and profile (which backend it runs on). Pass role to pick "
|
||||
+ "the contract — 'dev' implements a unit and opens its own PR, 'hunter' sweeps a "
|
||||
+ "scope without changing it, 'reviewer' reviews a diff it did not write, 'architect' "
|
||||
+ "refines a ticket before anyone builds it; omit "
|
||||
+ "the contract — 'dev' implements a unit and opens its own PR, 'reviewer' reviews a "
|
||||
+ "diff it did not write, 'architect' refines a ticket before anyone builds it; omit "
|
||||
+ "it for 'dev'. Pass profile (from fleet_profiles) to pick the backend, or omit it "
|
||||
+ "for the default. The two are independent: a reviewer may run on the same profile "
|
||||
+ "as the dev it reviews. The member opens your current directory by default; pass "
|
||||
@@ -2153,7 +2152,7 @@ public final class FleetMcp {
|
||||
+ "one. Returns the member's sessionId (use with fleet_send) and paneId (use with "
|
||||
+ "fleet_stop).",
|
||||
objectSchema(Map.of(
|
||||
"role", stringProp("What the member is for: architect, dev, hunter, or reviewer (default dev)"),
|
||||
"role", stringProp("What the member is for: architect, dev or reviewer (default dev)"),
|
||||
"profile", stringProp("Which backend to run it on (omit for the default profile)"),
|
||||
"cwd", stringProp("Working directory for the member (omit to inherit yours)"),
|
||||
"worktree", Map.of("type", "string", "description", "'true' or a ticket slug — requests an isolated git worktree"),
|
||||
|
||||
@@ -29,8 +29,8 @@ public enum MemberRole {
|
||||
* <p>Reads the repo and writes analysis. Never commits code and never opens a pull request —
|
||||
* an architect that starts implementing has stopped doing the job that makes it useful.
|
||||
*
|
||||
* <p>Architects are the one member kind with live slot binding, because a lead addresses the
|
||||
* same slots across many tickets and needs a stable name for them.
|
||||
* <p>Architects are the one member kind declared in config, because a lead addresses the same
|
||||
* slots across many tickets and needs a stable name for them.
|
||||
*/
|
||||
ARCHITECT,
|
||||
|
||||
@@ -43,14 +43,6 @@ public enum MemberRole {
|
||||
*/
|
||||
DEV,
|
||||
|
||||
/**
|
||||
* Sweeps an assigned package for defects and reports several ranked findings.
|
||||
*
|
||||
* <p>Never changes code, commits, or opens a pull request. A hunt gathers evidence, which can
|
||||
* include running the build, but leaves every fix to a later implementation unit.
|
||||
*/
|
||||
HUNTER,
|
||||
|
||||
/**
|
||||
* Reviews a diff it did not write and reports one structured finding.
|
||||
*
|
||||
@@ -67,7 +59,7 @@ public enum MemberRole {
|
||||
|
||||
/**
|
||||
* The {@code fleet:} block that holds this role's pool — {@code architects},
|
||||
* {@code developers}, {@code hunters}, {@code reviewers}.
|
||||
* {@code developers}, {@code reviewers}.
|
||||
*
|
||||
* <p>Plural, and not always the wire name: the pool of things a {@code dev} may run on reads
|
||||
* naturally as {@code developers:}. The wire name stays the singular {@code dev}, because that
|
||||
@@ -77,7 +69,6 @@ public enum MemberRole {
|
||||
return switch (this) {
|
||||
case ARCHITECT -> "architects";
|
||||
case DEV -> "developers";
|
||||
case HUNTER -> "hunters";
|
||||
case REVIEWER -> "reviewers";
|
||||
};
|
||||
}
|
||||
|
||||
@@ -0,0 +1,80 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.function.BiConsumer;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :591}. {@code Fleetd.main} wires {@link
|
||||
* dev.ltms.fleet.health.FleetHealthMonitor}'s {@code failTarget} callback with {@code
|
||||
* messages::abandon} — before this ticket that was an inline argument to {@code new
|
||||
* FleetHealthMonitor(...)}. Measured: replacing it with a no-op {@code BiConsumer} at the call
|
||||
* site compiles with 0 errors and leaves the full suite green, because nothing else in the tree
|
||||
* ever drives that specific constructor argument. In production it means a member the monitor
|
||||
* classifies {@code GONE}/{@code NEVER_READY} never has its pending ticket failed — the caller
|
||||
* keeps reporting {@code PENDING} for the full 30-minute async timeout instead of the immediate,
|
||||
* accurate failure CB-580 exists to give it.
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#healthFailTarget} directly — never {@code FleetHealthMonitor}
|
||||
* or {@code main} — against a real {@link MessageService}, using the same {@code sendAsync} +
|
||||
* {@code poll} observable {@link MessageServiceTest} already relies on to pin {@code
|
||||
* MessageService.abandon} itself.
|
||||
*/
|
||||
class FleetdHealthFailTargetWiringTest {
|
||||
|
||||
private static final String T = "term_a";
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.healthFailTarget delegates to the real MessageService.abandon, not a no-op")
|
||||
void healthFailTargetDelegatesToMessagesAbandon() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("$ prompt");
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
Injector injector = new Injector(agents);
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous);
|
||||
|
||||
BiConsumer<String, String> failTarget = Fleetd.healthFailTarget(messages);
|
||||
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
awaitWaiting(rendezvous);
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase());
|
||||
|
||||
failTarget.accept(T, "member unreachable (health monitor)");
|
||||
|
||||
MessageService.TaskView view = null;
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (System.currentTimeMillis() < deadline) {
|
||||
view = messages.poll(ticket);
|
||||
if (view.phase() != MessageService.Phase.PENDING) {
|
||||
break;
|
||||
}
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(10);
|
||||
}
|
||||
assertEquals(MessageService.Phase.FAILED, view == null ? null : view.phase(),
|
||||
"Fleetd.healthFailTarget(messages) must return messages::abandon — replacing it "
|
||||
+ "with a no-op BiConsumer at the Fleetd.healthFailTarget call site means "
|
||||
+ "this ticket is never failed and keeps polling as PENDING");
|
||||
assertTrue(view.detail() != null && view.detail().contains("member unreachable"),
|
||||
"the failure reason passed to failTarget.accept must reach MessageService.abandon "
|
||||
+ "and end up in the ticket's detail");
|
||||
}
|
||||
|
||||
private static void awaitWaiting(Rendezvous rendezvous) throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,48 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.net.ServerSocket;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :518}. {@code Fleetd.main} passes {@code LeadMailbox::open} as
|
||||
* the {@link Fleetd.LeadMailboxOpener} argument to {@code openLeadMailbox(...)} — before this
|
||||
* ticket that method reference was inline at the call site. Measured: replacing it with the inert
|
||||
* {@code (uri, selfCoordId, prefetch) -> null} compiles with 0 errors and leaves the full suite
|
||||
* green — {@code FleetdLeadMailboxSelectionTest} drives {@code openLeadMailbox} with its own
|
||||
* injected opener and never observes what {@code main} itself actually passes.
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#leadMailboxOpener} directly and proves it is the real,
|
||||
* network-attempting opener rather than a stub: pointed at a guaranteed-closed local port, it must
|
||||
* throw, exactly mirroring {@link FleetdReplyInboxOpenerWiringTest} for the reply-inbox opener.
|
||||
* The inert form never attempts a connection and returns {@code null} without throwing, so it
|
||||
* fails this assertion silently.
|
||||
*/
|
||||
class FleetdLeadMailboxOpenerWiringTest {
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.leadMailboxOpener is the real LeadMailbox::open, not a stub that never connects")
|
||||
void leadMailboxOpenerAttemptsARealConnection() throws Exception {
|
||||
int closedPort;
|
||||
try (ServerSocket socket = new ServerSocket(0)) {
|
||||
closedPort = socket.getLocalPort();
|
||||
} // released immediately — connecting to it now is a guaranteed refusal, not a fluke
|
||||
|
||||
Fleetd.LeadMailboxOpener opener = Fleetd.leadMailboxOpener();
|
||||
|
||||
IllegalStateException thrown = assertThrows(IllegalStateException.class,
|
||||
() -> opener.open("amqp://user:pw@127.0.0.1:" + closedPort + "/coord", "coord-1", 50),
|
||||
"Fleetd.leadMailboxOpener() must be LeadMailbox::open — a real network attempt "
|
||||
+ "against a genuinely unreachable broker must throw. The inert form "
|
||||
+ "(uri, selfCoordId, prefetch) -> null never attempts a connection and "
|
||||
+ "returns null instead of throwing, so it would fail this assertion "
|
||||
+ "silently.");
|
||||
assertTrue(thrown.getMessage().contains("cannot connect to AMQP coordination broker"),
|
||||
"must be LeadMailbox.open's own real failure message, not a different exception "
|
||||
+ "shape standing in for it");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,107 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.inject.Injector;
|
||||
import dev.ltms.fleet.mcp.PrimaryRegistry;
|
||||
import dev.ltms.fleet.msg.InMemoryReplyInbox;
|
||||
import dev.ltms.fleet.msg.MessageService;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.session.SessionManager;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.function.Consumer;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :611-631}. {@code Fleetd.main} wires {@code
|
||||
* sessions.onRelease(...)} with a lambda that calls three collaborators — {@code
|
||||
* messages.abandon}, {@code replyInbox.release}, and {@code primaryRegistry.forgetDelegation} —
|
||||
* before this ticket built inline inside {@code main}. Measured: replacing the whole lambda body
|
||||
* with {@code detail -> { }} compiles with 0 errors and leaves the full suite green, because each
|
||||
* collaborator is separately tested in isolation ({@code MessageServiceTest}, {@code
|
||||
* PrimaryRegistryTest}) but nothing before this ticket drove the lambda that calls all three from
|
||||
* {@code main}.
|
||||
*
|
||||
* <p>{@code MessageService.abandon}'s own javadoc documents the consequence: without this,
|
||||
* tearing a worker down leaves its rendezvous waiter open, so a blocking {@code fleet_send} keeps
|
||||
* blocking and an async one reports {@code PENDING} for a hardcoded thirty minutes on every {@code
|
||||
* fleet_stop} and every idle-reap.
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#releaseCleanup} directly — never {@code SessionManager} or
|
||||
* {@code main} — against real {@link MessageService}, {@link InMemoryReplyInbox}, and {@link
|
||||
* PrimaryRegistry} instances, and asserts each collaborator's own observable effect: the pending
|
||||
* async ticket transitions to {@code FAILED} (abandon), the inbox no longer owns the target's
|
||||
* queue (release), and the recorded delegation is forgotten (forgetDelegation).
|
||||
*/
|
||||
class FleetdReleaseCleanupWiringTest {
|
||||
|
||||
private static final String T = "term_a";
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.releaseCleanup reaches messages.abandon, replyInbox.release, and primaryRegistry.forgetDelegation")
|
||||
void releaseCleanupReachesAllThreeCollaborators() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("$ prompt");
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
Injector injector = new Injector(agents);
|
||||
MessageService messages = new MessageService(agents, injector, rendezvous);
|
||||
InMemoryReplyInbox replyInbox = new InMemoryReplyInbox();
|
||||
PrimaryRegistry primaryRegistry = new PrimaryRegistry(null);
|
||||
|
||||
// Set up the "before" state each collaborator's own effect is measured against.
|
||||
String ticket = messages.sendAsync(T, "long task");
|
||||
awaitWaiting(rendezvous);
|
||||
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase(),
|
||||
"sanity: the async ticket is pending before cleanup runs");
|
||||
|
||||
replyInbox.own(T);
|
||||
replyInbox.publish(T, "msg-1", "hello");
|
||||
assertEquals(1, replyInbox.peek(T).size(),
|
||||
"sanity: the inbox owns T and holds one message before cleanup runs");
|
||||
|
||||
primaryRegistry.recordDelegation(T, "lead-1");
|
||||
assertEquals("lead-1", primaryRegistry.nudgeTargetFor(T).orElse(null),
|
||||
"sanity: the delegation is recorded before cleanup runs");
|
||||
|
||||
Consumer<SessionManager.ReleaseDetail> cleanup =
|
||||
Fleetd.releaseCleanup(messages, replyInbox, primaryRegistry);
|
||||
cleanup.accept(new SessionManager.ReleaseDetail(T, null, null, null, null));
|
||||
|
||||
MessageService.TaskView view = null;
|
||||
long deadline = System.currentTimeMillis() + 3000;
|
||||
while (System.currentTimeMillis() < deadline) {
|
||||
view = messages.poll(ticket);
|
||||
if (view.phase() != MessageService.Phase.PENDING) {
|
||||
break;
|
||||
}
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(10);
|
||||
}
|
||||
assertEquals(MessageService.Phase.FAILED, view == null ? null : view.phase(),
|
||||
"releaseCleanup must call messages.abandon(...) — an inert detail -> { } lambda "
|
||||
+ "leaves this ticket PENDING forever");
|
||||
assertTrue(view.detail() != null && view.detail().contains("released"),
|
||||
"the abandon reason must say the worker session was released");
|
||||
|
||||
assertTrue(replyInbox.peek(T).isEmpty(),
|
||||
"releaseCleanup must call replyInbox.release(...) — an inert lambda leaves the "
|
||||
+ "inbox still owning T with its message");
|
||||
|
||||
assertTrue(primaryRegistry.nudgeTargetFor(T).isEmpty(),
|
||||
"releaseCleanup must call primaryRegistry.forgetDelegation(...) — an inert lambda "
|
||||
+ "leaves the stale delegation in place");
|
||||
}
|
||||
|
||||
private static void awaitWaiting(Rendezvous rendezvous) throws InterruptedException {
|
||||
long deadline = System.currentTimeMillis() + 2000;
|
||||
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||
//noinspection BusyWait
|
||||
Thread.sleep(5);
|
||||
}
|
||||
assertTrue(rendezvous.isWaiting(T), "send should have opened its rendezvous waiter");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.net.ServerSocket;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :512}. {@code Fleetd.main} passes {@code AmqpReplyInbox::open}
|
||||
* as the {@link Fleetd.AmqpOpener} argument to {@code selectReplyInbox(...)} — before this ticket
|
||||
* that method reference was inline at the call site. Measured: replacing it with the inert {@code
|
||||
* (uri, prefetch) -> new InMemoryReplyInbox()} compiles with 0 errors and leaves the full suite
|
||||
* green — {@code FleetdReplyInboxSelectionTest} drives {@code selectReplyInbox} with its own
|
||||
* injected opener (including one test that passes the real {@code AmqpReplyInbox::open}
|
||||
* explicitly) and never observes what {@code main} itself actually passes.
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#replyInboxOpener} directly and proves it is the real, network-
|
||||
* attempting opener rather than a stub: pointed at a guaranteed-closed local port, it must throw —
|
||||
* the same shape {@code FleetdReplyInboxSelectionTest.aRealUnreachableBrokerFallsBackViaTheRealOpener}
|
||||
* already relies on for {@code AmqpReplyInbox.open} itself. The inert form never attempts a
|
||||
* connection and never throws, so it fails this assertion silently (by returning normally).
|
||||
*/
|
||||
class FleetdReplyInboxOpenerWiringTest {
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.replyInboxOpener is the real AmqpReplyInbox::open, not a stub that never connects")
|
||||
void replyInboxOpenerAttemptsARealConnection() throws Exception {
|
||||
int closedPort;
|
||||
try (ServerSocket socket = new ServerSocket(0)) {
|
||||
closedPort = socket.getLocalPort();
|
||||
} // released immediately — connecting to it now is a guaranteed refusal, not a fluke
|
||||
|
||||
Fleetd.AmqpOpener opener = Fleetd.replyInboxOpener();
|
||||
|
||||
IllegalStateException thrown = assertThrows(IllegalStateException.class,
|
||||
() -> opener.open("amqp://user:pw@127.0.0.1:" + closedPort + "/vh", 50),
|
||||
"Fleetd.replyInboxOpener() must be AmqpReplyInbox::open — a real network attempt "
|
||||
+ "against a genuinely unreachable broker must throw. The inert form "
|
||||
+ "(uri, prefetch) -> new InMemoryReplyInbox() never attempts a connection "
|
||||
+ "and never throws, so it would return normally here and fail this "
|
||||
+ "assertion silently.");
|
||||
assertTrue(thrown.getMessage().contains("cannot connect to AMQP broker"),
|
||||
"must be AmqpReplyInbox.open's own real failure message, not a different exception "
|
||||
+ "shape standing in for it");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,72 @@
|
||||
package dev.ltms.fleet;
|
||||
|
||||
import dev.ltms.fleet.herdr.AgentControl;
|
||||
import dev.ltms.fleet.herdr.FakeHerdr;
|
||||
import dev.ltms.fleet.inject.CompletionResolver;
|
||||
import dev.ltms.fleet.inject.ExhaustedPatternLookup;
|
||||
import dev.ltms.fleet.inject.ExhaustionSink;
|
||||
import dev.ltms.fleet.inject.TurnRegistrar;
|
||||
import dev.ltms.fleet.msg.Rendezvous;
|
||||
import dev.ltms.fleet.msg.TurnToken;
|
||||
import org.junit.jupiter.api.DisplayName;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* fleetd #589 Group 3, site {@code :505}. {@code Fleetd.main} wires the {@link
|
||||
* dev.ltms.fleet.inject.Injector}'s {@link TurnRegistrar} with {@code completion::register} —
|
||||
* before this ticket that was an inline argument to {@code new Injector(...)}, so nothing could
|
||||
* pin it directly. Measured: replacing it with {@link TurnRegistrar#NOOP} at the call site
|
||||
* compiles with 0 errors and leaves the full suite green, because {@code onDelivered}'s own {@code
|
||||
* captureBaseline} performs the identical {@code inFlight} check-and-put a moment later on the
|
||||
* ordinary delivery path — the two are indistinguishable unless something reads the resolver
|
||||
* between {@code register} and {@code onDelivered}, or {@code onDelivered} never runs at all (the
|
||||
* gap fleetd #556 introduced {@link TurnRegistrar} to close).
|
||||
*
|
||||
* <p>This test calls {@link Fleetd#turnRegistrar} directly — never {@code Injector} or {@code
|
||||
* main} — and drives {@link CompletionResolver} entirely through its public API: {@link
|
||||
* TurnRegistrar#register} followed by {@link CompletionResolver#resolveBeforePostAction}, which
|
||||
* looks up the same {@code inFlight} entry {@code onTurnComplete} would. With the real registrar,
|
||||
* that entry exists and the waiter opened by {@link Rendezvous#open} resolves; with {@link
|
||||
* TurnRegistrar#NOOP} nothing was ever registered, {@code resolveBeforePostAction} finds no
|
||||
* in-flight turn, and the waiter is left exactly as it started — never done.
|
||||
*/
|
||||
class FleetdTurnRegistrarWiringTest {
|
||||
|
||||
private static final String T = "term_a";
|
||||
|
||||
@Test
|
||||
@DisplayName("Fleetd.turnRegistrar delegates to the real CompletionResolver, not a no-op")
|
||||
void turnRegistrarDelegatesToCompletionRegister() {
|
||||
FakeHerdr herdr = new FakeHerdr().readText("⏺ BUILD GREEN: 391 files\n❯ ");
|
||||
AgentControl agents = new AgentControl(herdr);
|
||||
Rendezvous rendezvous = new Rendezvous();
|
||||
// fleetd#164: an ever-advancing fake clock stands in for the real time a turn would take
|
||||
// between delivery and resolution, so the MIN_TURN_NANOS "too fast" floor never trips here —
|
||||
// see MessageServiceTest's resolverClock for the same technique.
|
||||
AtomicLong clock = new AtomicLong();
|
||||
CompletionResolver completion = new CompletionResolver(agents, rendezvous,
|
||||
ExhaustedPatternLookup.none(), ExhaustionSink.none(),
|
||||
() -> clock.addAndGet(CompletionResolver.MIN_TURN_NANOS + 1));
|
||||
|
||||
TurnRegistrar registrar = Fleetd.turnRegistrar(completion);
|
||||
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open(T);
|
||||
registrar.register(T, new TurnToken(T, waiter));
|
||||
|
||||
// Mirrors what Injector.onStatus's confirmed working->idle boundary would trigger via
|
||||
// CompletionResolver.onTurnComplete — resolveBeforePostAction is the public, synchronous
|
||||
// twin of that path and reads the exact same inFlight entry register() must have written.
|
||||
completion.resolveBeforePostAction(T);
|
||||
|
||||
assertTrue(waiter.isDone(),
|
||||
"Fleetd.turnRegistrar(completion) must return completion::register — replacing it "
|
||||
+ "with TurnRegistrar.NOOP at the Fleetd.turnRegistrar call site means this "
|
||||
+ "turn is never registered with CompletionResolver, so resolveBeforePostAction "
|
||||
+ "finds no in-flight turn and this waiter is never resolved");
|
||||
}
|
||||
}
|
||||
@@ -345,7 +345,7 @@ class FleetConfigTest {
|
||||
IllegalStateException unknownError = assertThrows(IllegalStateException.class,
|
||||
() -> FleetConfig.load(unknown).validateCharters());
|
||||
assertTrue(unknownError.getMessage().contains("architetc"));
|
||||
assertTrue(unknownError.getMessage().contains("[architect, dev, hunter, reviewer]"));
|
||||
assertTrue(unknownError.getMessage().contains("[architect, dev, reviewer]"));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -812,17 +812,13 @@ class FleetConfigTest {
|
||||
reviewers:
|
||||
b:
|
||||
profile: sonnet
|
||||
hunters:
|
||||
c:
|
||||
profile: sonnet
|
||||
""");
|
||||
|
||||
FleetConfig cfg = FleetConfig.load(f);
|
||||
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.DEV));
|
||||
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.HUNTER));
|
||||
assertEquals(List.of("sonnet"), cfg.fleet().profilesFor(MemberRole.REVIEWER));
|
||||
assertTrue(cfg.fleet().profilesFor(MemberRole.ARCHITECT).isEmpty());
|
||||
assertEquals(List.of(MemberRole.DEV, MemberRole.HUNTER, MemberRole.REVIEWER), cfg.fleet().rolesConfigured());
|
||||
assertEquals(List.of(MemberRole.DEV, MemberRole.REVIEWER), cfg.fleet().rolesConfigured());
|
||||
}
|
||||
|
||||
/** The case the two axes exist for: one backend, two roles, and neither is a duplicate. */
|
||||
|
||||
@@ -1937,7 +1937,7 @@ class FleetMcpTest {
|
||||
null, null, null, null, null, null);
|
||||
|
||||
assertEquals(Boolean.TRUE, res.isError());
|
||||
assertTrue(textOf(res).contains("architect, dev, hunter, reviewer"), textOf(res));
|
||||
assertTrue(textOf(res).contains("architect, dev, reviewer"), textOf(res));
|
||||
}
|
||||
|
||||
// ── CB-619 / fleetd #123: a spawn asking for a role its profile has no slot for must be
|
||||
|
||||
@@ -556,31 +556,6 @@ class ClaudeCodeLauncherTest {
|
||||
"no --agent flag when the role has no agent-definition file");
|
||||
}
|
||||
|
||||
@Test
|
||||
void hunterRoleUsesItsAgentFileAndStopsUsingItWhenRemoved(@TempDir Path cwd) throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
Path agentFile = Files.createDirectories(cwd.resolve(".claude/agents")).resolve("hunter.md");
|
||||
Files.writeString(agentFile, "---\nname: hunter\n---\nSweep for defects.");
|
||||
FleetConfig.Profile cfg = new FleetConfig.Profile(
|
||||
"sonnet", "http://gx00.gw:8000", null, null, "FLEETD_WORKER_TOKEN",
|
||||
List.of("claude"), "tab", "fleetd-workers", "w #{n}", null, null, null);
|
||||
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||
|
||||
svc.spawn(new SpawnRequest("sonnet", cwd.toString(), null, null, null, MemberRole.HUNTER));
|
||||
|
||||
List<String> args = spawnedArgs(herdr);
|
||||
int flag = args.indexOf("--agent");
|
||||
assertTrue(flag >= 0, "the hunter role reaches its agent-definition file: " + args);
|
||||
assertEquals("hunter", args.get(flag + 1));
|
||||
|
||||
Files.delete(agentFile);
|
||||
svc.spawn(new SpawnRequest("sonnet", cwd.toString(), null, null, null, MemberRole.HUNTER));
|
||||
|
||||
assertFalse(spawnedArgs(herdr).contains("--agent"),
|
||||
"the hunter role no longer gets an agent when its file is removed");
|
||||
}
|
||||
|
||||
private ClaudeCodeLauncher multiProfile(FakeHerdr herdr) {
|
||||
FleetConfig.Profile gx10 = new FleetConfig.Profile("gx10", "http://gx10.gw:8000", "coder",
|
||||
null, "FLEETD_WORKER_TOKEN", List.of("claude"), "tab", "fleetd-workers", "w #{n}", null, null, null);
|
||||
|
||||
@@ -10,13 +10,12 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
class MemberRoleTest {
|
||||
|
||||
@Test
|
||||
void theFourRolesAreArchitectDevHunterAndReviewer() {
|
||||
assertEquals(4, MemberRole.values().length,
|
||||
void theThreeRolesAreArchitectDevAndReviewer() {
|
||||
assertEquals(3, MemberRole.values().length,
|
||||
"a new role changes the charter, the role file, the skill and the authz row — "
|
||||
+ "adding one is a deliberate act, so this count is meant to fail first");
|
||||
assertEquals("architect", MemberRole.ARCHITECT.wireName());
|
||||
assertEquals("dev", MemberRole.DEV.wireName());
|
||||
assertEquals("hunter", MemberRole.HUNTER.wireName());
|
||||
assertEquals("reviewer", MemberRole.REVIEWER.wireName());
|
||||
}
|
||||
|
||||
@@ -31,7 +30,6 @@ class MemberRoleTest {
|
||||
void parseIsCaseInsensitiveAndTrimsSurroundingSpace() {
|
||||
assertSame(MemberRole.ARCHITECT, MemberRole.parse("Architect"));
|
||||
assertSame(MemberRole.DEV, MemberRole.parse(" DEV "));
|
||||
assertSame(MemberRole.HUNTER, MemberRole.parse("HuNtEr"));
|
||||
assertSame(MemberRole.REVIEWER, MemberRole.parse("ReViEwEr"));
|
||||
}
|
||||
|
||||
@@ -40,7 +38,7 @@ class MemberRoleTest {
|
||||
IllegalArgumentException e =
|
||||
assertThrows(IllegalArgumentException.class, () -> MemberRole.parse("archtiect"));
|
||||
assertTrue(e.getMessage().contains("archtiect"), e.getMessage());
|
||||
assertTrue(e.getMessage().contains("architect, dev, hunter, reviewer"),
|
||||
assertTrue(e.getMessage().contains("architect, dev, reviewer"),
|
||||
"a typo in config should be fixable from the message alone: " + e.getMessage());
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user