Compare commits
4 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| de026b8f8a | |||
| 97f6c33a45 | |||
| 7d4a4339c2 | |||
| ee5f8b932b |
+166
@@ -0,0 +1,166 @@
|
||||
# CB-137 / fleetd issue #137 — report
|
||||
|
||||
## Real root cause (not the hypothesis in the ticket)
|
||||
|
||||
I read `MessageService.java` and `Rendezvous.java` before changing anything. The mechanism is real,
|
||||
but the exact place it happens is `MessageService.answer()`, not "the reply goes to the inbox on
|
||||
purpose" in general.
|
||||
|
||||
1. A lead delegates with `fleet_send{wait:false}` → `sendAsync()` creates a `Task` and runs `send()`
|
||||
on a background virtual thread with a 30-minute internal budget (`ASYNC_TIMEOUT_MS`).
|
||||
2. The worker calls `fleet_ask`. That resolves the open rendezvous waiter with `Kind.QUESTION`, so
|
||||
`send()` returns immediately and the `Task` is left open (its `future` stays unresolved — see the
|
||||
comment in `sendAsync`'s lambda: "Keep the accepted owner until answer() finishes it").
|
||||
3. The lead answers with `fleet_send{turnId, content}`. This calls `FleetMcp.answer()` →
|
||||
`MessageService.answer(turnId, content, timeout)`. The `timeout` here is **not** the generous
|
||||
30-minute async budget — it is the MCP tool's own bounded wait: `DEFAULT_TIMEOUT_MS = 25_000`,
|
||||
clamped to at most `MAX_TIMEOUT_MS = 120_000` (`FleetMcp.java:71-72,495,512`). This is the same
|
||||
~60–120s window every blocking `fleet_send` call is capped at (documented elsewhere as "the
|
||||
caller's own MCP client call timeout").
|
||||
4. `answer()` opens a **fresh** rendezvous waiter for the worker session and blocks on it for at most
|
||||
that window. If the worker's resumed turn takes longer than that to actually finish (very
|
||||
plausible — the resumed turn can mean more edits, a build, a commit, a push, opening a PR), the
|
||||
wait times out. On timeout, `answer()`'s `finally` block unconditionally calls
|
||||
`rendezvous.close(workerSession, reply)`, **removing the waiter from the map**, and returns
|
||||
`Outcome.TIMED_OUT_WORKING` to the lead.
|
||||
5. The worker keeps working, unaware anything happened, and eventually calls `fleet_reply`. That
|
||||
reaches `MessageService.reply(session, content)`, which tries `rendezvous.resolve(session,
|
||||
content)` — but the waiter was already closed in step 4, so `resolve` returns `false`. `reply()`
|
||||
then falls back to `inbox.publish(...)` and marks `strandedReplies.put(session, true)`
|
||||
(CB-640 bookkeeping) — the reply is safely held, but **the async `Task`'s `future` is never
|
||||
completed**.
|
||||
6. `fleet_poll{ticket}` keeps returning `PENDING` forever (the `Task` never resolves) — until the
|
||||
lead eventually calls `fleet_stop`. That fires `sessions.onRelease` → `messages.abandon(target,
|
||||
reason)` (`Fleetd.java:481-497`), where `reason` is built with the exact text from the bug report
|
||||
("the worker session was released before it replied; worktree=... branch=... snapshot=...",
|
||||
`Fleetd.java:484-487`). `abandon()`'s loop finds the still-open `Task` (`question == null`, future
|
||||
not done) and completes it as `WORKER_FAILED` with that misleading reason — even though the
|
||||
worker's real reply is sitting, intact, in the inbox the whole time.
|
||||
|
||||
So: the reported behaviour is correct, and the specific trigger is `answer()`'s own bounded wait
|
||||
being shorter than the worker's real resumed-turn time — not anything to do with the ~55s
|
||||
`fleet_ask` window itself (that part, issue #61, is untouched).
|
||||
|
||||
## Fix
|
||||
|
||||
Two changes in `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`, both scoped to the
|
||||
ticket/reply routing and the terminal-state text — `fleet_ask`'s own window and mechanics are
|
||||
untouched.
|
||||
|
||||
**1. `reply()` — priority 1 (the ticket resolves with the real reply).**
|
||||
Before falling back to the inbox, `reply()` now looks for an async `Task` that is specifically in the
|
||||
"already answered but not yet resolved" state (`question == null`, `turnId != null` — set once
|
||||
`answer()` has cleared the question but before anything completed the future, `!future.isDone()`).
|
||||
If one exists for this `target`, the worker's reply completes that `Task`'s future directly as
|
||||
`Outcome.REPLIED` with the real content, and the reply never touches the inbox at all. A task that
|
||||
was never asked has `turnId == null` and can never match, so ordinary (no-`fleet_ask`) delegations
|
||||
are unaffected — they already resolve through the pre-existing rendezvous fast path.
|
||||
|
||||
I chose this over leaving `answer()`'s own timeout behaviour untouched and instead keeping its
|
||||
rendezvous waiter open in the background: that alternative works but reopens the "at most one
|
||||
waiter per session" invariant (`Rendezvous.open` throws on a double-open) to a new class of races
|
||||
with a fresh send arriving mid-window. The `send()` path already guards against sending into an
|
||||
answered-but-still-resolving worker via `hasAsyncQuestion(target)` (checks `asyncTasksByTurn`,
|
||||
which still holds the task until it resolves), so routing through `reply()` gets the same protection
|
||||
without touching `answer()`'s waiter lifecycle at all — the smaller, safer diff.
|
||||
|
||||
**2. `abandon()` — priority 3 (required independently, "even if you fix (1)").**
|
||||
Before marking any of a released target's still-open tasks `WORKER_FAILED`, `abandon()` now checks
|
||||
`hasStrandedReply(target)` (the existing CB-640 fact — true whenever the *last* `reply()` for this
|
||||
target fell through to the inbox). If true, it drains the inbox (`recoverStrandedReply`) and — if it
|
||||
actually finds a message — completes the task as `REPLIED` with that real content instead of writing
|
||||
the failure. This is deliberately a **separate** check from fix 1: fix 1 already prevents the
|
||||
inbox-stranding from happening in the exact scenario this ticket describes, so by the time
|
||||
`abandon()` runs the task is normally already resolved and `abandon()`'s `complete()` call is a
|
||||
harmless no-op. This second check exists so that if some *other* future path ever strands a reply
|
||||
in the inbox without resolving its ticket, `abandon()` still refuses to report a false failure —
|
||||
"if a reply reached any sink for that turn, the terminal state is done," per the ticket. I verified
|
||||
both are required by disabling each independently and confirming the two new tests fail (see below).
|
||||
|
||||
**Priority 4 (the snapshot/worktree hint).** Handled as a consequence of both fixes rather than a
|
||||
separate branch: once a task resolves as `REPLIED` (via either fix), `abandon()` never calls
|
||||
`new Reply(Outcome.WORKER_FAILED, reason)` for that task at all, so the "the worker session was
|
||||
released before it replied; worktree=... branch=... snapshot=..." text is never constructed or
|
||||
attached to that ticket's outcome. It still appears, correctly, for a task that never got a reply
|
||||
(the existing `abandonFailsEveryPendingAsyncTicketForTheReleasedTarget` /
|
||||
`anAbandonedAsyncTaskPollsAsFailedNotPending` tests still pass unchanged).
|
||||
|
||||
**Priority 2** was not needed — fix 1 makes `fleet_poll{ticket}` return the actual reply (the
|
||||
higher-priority option), so I did not fall back to "the ticket merely resolves as done with no
|
||||
content."
|
||||
|
||||
## Tests — driven through the real delegation path, not the reply sink directly
|
||||
|
||||
Both new tests in `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java` go through
|
||||
`sendAsync` → `injectDelivery` → `ask` → `answer` (with a short timeout, so it genuinely times out,
|
||||
mirroring the ~25–120s real MCP-call bound vs. a longer resumed turn) → `reply` → `poll`/`abandon`.
|
||||
No test constructs a `Reply` and hands it to a sink directly.
|
||||
|
||||
- `aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket` — asserts `fleet_poll{ticket}` (via
|
||||
`messages.poll`) reaches `Phase.DONE` with the worker's actual reply text and
|
||||
`replySource() == "reply"`, and that `hasStrandedReply(T)` stays `false` (proves the reply never
|
||||
touched the inbox at all — fix 1 caught it).
|
||||
- `fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket` — same setup, then calls `abandon(T, "the
|
||||
worker session was released before it replied")` (what `fleet_stop` triggers) and asserts it
|
||||
returns `false` (no failure recorded) and the ticket still polls `DONE` with the real reply.
|
||||
|
||||
**Proof both fail without the change.** I temporarily short-circuited both new private methods
|
||||
(`askAnsweredAsyncTask` → always `null`, `recoverStrandedReply` → always `null`) — i.e. disabled
|
||||
both fixes — and ran just these two tests:
|
||||
|
||||
```
|
||||
[ERROR] Tests run: 2, Failures: 2, Errors: 0, Skipped: 0
|
||||
dev.ltms.fleet.msg.MessageServiceTest.aReplyAfterAnswerTimesOutStillCompletesTheAsyncTicket
|
||||
org.opentest4j.AssertionFailedError: expected: <DONE> but was: <PENDING>
|
||||
dev.ltms.fleet.msg.MessageServiceTest.fleetStopAfterAnOrphanedReplyDoesNotFailTheTicket
|
||||
org.opentest4j.AssertionFailedError: a reply already arrived, so nothing here is a genuine failure
|
||||
==> expected: <false> but was: <true>
|
||||
```
|
||||
|
||||
This is the exact bug: the ticket stays `PENDING` forever, and `abandon()` reports `true` (a
|
||||
failure) even though a reply had already arrived. I then restored both fixes (verified with
|
||||
`grep -n "TEMP #137-proof"` finding nothing) and re-ran — both pass.
|
||||
|
||||
## Build
|
||||
|
||||
Ran from `fleetd/`, unpiped, full output read (not `| tail`):
|
||||
|
||||
```
|
||||
mvn clean install
|
||||
...
|
||||
[INFO] Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0
|
||||
[INFO] BUILD SUCCESS
|
||||
[INFO] Total time: 36.315 s
|
||||
```
|
||||
|
||||
Main was at 1037 tests; this branch adds the 2 new tests above → 1039, all green, `exit=0`.
|
||||
|
||||
## What I could NOT check
|
||||
|
||||
- No IDE tooling is mounted for me (worker), so no `ide_diagnostics`/IntelliJ inspection pass — only
|
||||
`mvn clean install` (compiler + full test suite), as the worker procedure allows.
|
||||
- I cannot restart the daemon or dogfood this live — I have no forge/daemon control. This is
|
||||
unverified against a real herdr pane, a real MCP client's ~60s call cap, or a real worker session;
|
||||
everything above is verified only through the JUnit fixture's simulated timing
|
||||
(`FakeHerdr`/`injector.onStatus`/direct `messages.answer(...,150)` calls), not a live fleet.
|
||||
A primary should still consider a short live dogfood (an async delegation that asks, gets answered,
|
||||
and takes longer than ~2 minutes to reply) before calling this closed.
|
||||
- I did not touch, and did not re-verify, the `fleet_ask` ~55s window itself (issue #61) — out of
|
||||
scope per the brief.
|
||||
|
||||
## Scope note (not investigated further)
|
||||
|
||||
`answer()`'s nested/double-`fleet_ask` case (the worker asks a second question before ever
|
||||
replying to the first answer) has some pre-existing behaviour around which `turnId` a `QUESTION`
|
||||
resolution gets attributed to that I did not fully untangle — it predates this change, my fix does
|
||||
not touch it, and it is unrelated to the reported defect. Flagging only; not investigated further.
|
||||
|
||||
## Handoff
|
||||
|
||||
- Branch: `worker/cb-137-ask-ticket-e7760c-2`
|
||||
- Worktree root: `/Users/dai.ha/LTMS/.bridged-worktrees/734324-2`
|
||||
- Files changed:
|
||||
- `fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java`
|
||||
- `fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java`
|
||||
- `REPORT-cb137.md` (this file)
|
||||
- Build: `Tests run: 1039, Failures: 0, Errors: 0, Skipped: 0` / `BUILD SUCCESS` (verbatim above)
|
||||
@@ -951,10 +951,7 @@ public final class FleetMcp {
|
||||
.sorted(Map.Entry.comparingByValue())
|
||||
.map(e -> leadView(e.getKey(), e.getValue(), live.get(e.getKey()), selfTerm))
|
||||
.toList();
|
||||
// fleetd #209: this is the caller-driven fleet_list read that actually reports
|
||||
// agentSessionId (via memberCapacityView -> SessionManager.rosterView), so it uses the
|
||||
// resolving roster; the heartbeat/health/metrics timers stay on the plain sessions.roster().
|
||||
List<MemberSession> roster = sessions.rosterResolved();
|
||||
List<MemberSession> roster = sessions.roster();
|
||||
List<Map<String, Object>> out = roster.stream()
|
||||
.map(s -> memberCapacityView(s, live.get(s.terminalId()), messages, capacity.clock().getAsLong()))
|
||||
.toList();
|
||||
|
||||
@@ -9,6 +9,7 @@ import dev.ltms.fleet.metrics.Metrics;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
@@ -167,6 +168,12 @@ public final class MessageService {
|
||||
private static final class Task {
|
||||
private final String ticket;
|
||||
private final String target;
|
||||
/**
|
||||
* When this task was created (#137 fix): the tiebreaker for which of several open tasks on
|
||||
* one target gets a recovered reply in {@link #abandon} — the oldest, since it is the one
|
||||
* that has been waiting longest.
|
||||
*/
|
||||
private final long createdNanos;
|
||||
private final CompletableFuture<Reply> future = new CompletableFuture<>();
|
||||
/**
|
||||
* When {@link #future} resolved, or {@code null} while it is still pending — the clock
|
||||
@@ -184,6 +191,7 @@ public final class MessageService {
|
||||
private Task(String ticket, String target, LongSupplier nowNanos) {
|
||||
this.ticket = ticket;
|
||||
this.target = target;
|
||||
this.createdNanos = nowNanos.getAsLong();
|
||||
future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong());
|
||||
}
|
||||
}
|
||||
@@ -375,6 +383,19 @@ public final class MessageService {
|
||||
* {@link Rendezvous#resolveQuestion} must keep today's {@code NO_WAITER} behaviour — questions
|
||||
* are interactive and must never be queued.
|
||||
*
|
||||
* <p><strong>Ambiguous match also falls to the inbox.</strong> {@link #askAnsweredAsyncTasks}
|
||||
* cannot actually return more than one entry today (see its own javadoc for why — in short,
|
||||
* {@link #hasAsyncQuestion} keeps a target BUSY, so no second task can reach this state, for as
|
||||
* long as an earlier one's {@code turnId} is still stamped). That is an emergent guarantee from
|
||||
* two other facts, not one this method enforces, so this branch stays in as defence in depth
|
||||
* rather than being removed as dead code: if it ever weakens, returning whichever candidate a
|
||||
* {@code ConcurrentHashMap} iteration reaches first would let a genuine reply complete the
|
||||
* <em>wrong</em> ticket — silently handing the lead something that reads like a correct answer to
|
||||
* a delegation the worker never touched, which is worse than a failure because the lead acts on
|
||||
* it. When more than one candidate exists, guessing is not safe: fall back to the inbox exactly
|
||||
* as the zero-candidate case does, and let {@link #abandon} apply the eventual recovery
|
||||
* deterministically instead.
|
||||
*
|
||||
* @return always {@code true} — the reply resolved a live send, completed a parked ticket, or
|
||||
* was queued
|
||||
*/
|
||||
@@ -392,13 +413,21 @@ public final class MessageService {
|
||||
// FAILED with a misleading "session released before it replied" reason, even though the reply
|
||||
// had, in fact, arrived. Completing the matching ticket directly here means fleet_poll{ticket}
|
||||
// sees the real reply instead.
|
||||
Task orphan = askAnsweredAsyncTask(session);
|
||||
if (orphan != null && orphan.future.complete(new Reply(Outcome.REPLIED, content))) {
|
||||
if (orphan.turnId != null) {
|
||||
asyncTasksByTurn.remove(orphan.turnId, orphan);
|
||||
List<Task> candidates = askAnsweredAsyncTasks(session);
|
||||
if (candidates.size() == 1) {
|
||||
Task orphan = candidates.get(0);
|
||||
if (orphan.future.complete(new Reply(Outcome.REPLIED, content))) {
|
||||
if (orphan.turnId != null) {
|
||||
asyncTasksByTurn.remove(orphan.turnId, orphan);
|
||||
}
|
||||
count(FleetMetrics.REPLIES, "path", "async-recovered");
|
||||
return true; // the ticket itself took it — no inbox stranding at all
|
||||
}
|
||||
count(FleetMetrics.REPLIES, "path", "async-recovered");
|
||||
return true; // the ticket itself took it — no inbox stranding at all
|
||||
} else if (candidates.size() > 1) {
|
||||
List<String> tickets = candidates.stream().map(t -> t.ticket).toList();
|
||||
log.warn("reply from {} matches {} open async tickets {} — cannot tell which one it "
|
||||
+ "answers, queuing to the inbox instead of guessing", session, candidates.size(),
|
||||
tickets);
|
||||
}
|
||||
inbox.publish(session, UUID.randomUUID().toString(), content);
|
||||
// CB-640: record the stranding itself (not just the reply text) so fleet health can see a
|
||||
@@ -414,21 +443,38 @@ public final class MessageService {
|
||||
}
|
||||
|
||||
/**
|
||||
* The still-open async task on {@code target} whose {@code fleet_ask} was already answered — its
|
||||
* {@link Task#turnId} is stamped but its {@link Task#question} was cleared by {@link #answer} —
|
||||
* yet whose future is not resolved yet (#137). {@code null} if no such task exists, including the
|
||||
* Every still-open async task on {@code target} whose {@code fleet_ask} was already answered —
|
||||
* its {@link Task#turnId} is stamped but its {@link Task#question} was cleared by {@link #answer}
|
||||
* — yet whose future is not resolved yet (#137). Empty if no such task exists, including the
|
||||
* common case where {@code target}'s worker never used {@code fleet_ask} at all (a task that was
|
||||
* never asked has {@code turnId == null}, so it can never match here and only ever completes
|
||||
* through the ordinary rendezvous fast path in {@link #reply}).
|
||||
*
|
||||
* <p><strong>Returns at most one entry today — verified, not assumed.</strong> {@link #send}
|
||||
* refuses to open a waiter on {@code target} while {@link #hasAsyncQuestion} is true, and that
|
||||
* check matches ANY task whose {@code turnId} is still stamped in {@code asyncTasksByTurn} —
|
||||
* not only while its question is still open. {@link #answer} deliberately leaves that stamp in
|
||||
* place ({@code clearAsyncQuestion(turnId, false)}) until the resumed turn's own future actually
|
||||
* resolves, at which point {@link #finishAsyncTask} both removes the stamp AND completes that
|
||||
* task's future in the same call. So a second task can never reach "{@code turnId} stamped, future
|
||||
* still open" — the exact pair this method matches on — while a first one already holds it: by
|
||||
* the time the stamp is gone, so is the eligibility. This is an emergent property of those two
|
||||
* facts holding together, not something this method (or its callers) enforces on its own — flip
|
||||
* {@code forgetTurn} to {@code true} in that one {@link #answer} call and it silently stops being
|
||||
* true, with nothing left to fail loudly. The callers below still handle "more than one" as
|
||||
* defence in depth against exactly that, not because they exercise it today: {@link #reply}
|
||||
* treats it as unresolvable and falls back to the inbox; {@link #abandon} would pick the oldest
|
||||
* deterministically (its own {@code matching} list has no such guarantee — see its javadoc).
|
||||
*/
|
||||
private Task askAnsweredAsyncTask(String target) {
|
||||
private List<Task> askAnsweredAsyncTasks(String target) {
|
||||
List<Task> candidates = new ArrayList<>();
|
||||
for (Task task : tasks.values()) {
|
||||
if (target.equals(task.target) && task.question == null && task.turnId != null
|
||||
&& !task.future.isDone()) {
|
||||
return task;
|
||||
candidates.add(task);
|
||||
}
|
||||
}
|
||||
return null;
|
||||
return candidates;
|
||||
}
|
||||
|
||||
/** Record a counter sample when a registry is wired; a no-op in unit tests. */
|
||||
@@ -473,16 +519,55 @@ public final class MessageService {
|
||||
* outcome is counted, so a torn-down delegation stops being invisible to {@code /metrics}.
|
||||
*
|
||||
* <p><strong>#137 defence in depth.</strong> {@link #reply} already hands a worker's real
|
||||
* {@code fleet_reply} straight to the async ticket it belongs to whenever one is still parked
|
||||
* waiting for it (see {@link #askAnsweredAsyncTask}), so by the time a session is released its
|
||||
* tasks are normally already resolved — this loop's {@code complete} calls are then harmless
|
||||
* no-ops (a {@link CompletableFuture} can only resolve once). But should some other path someday
|
||||
* strand a reply in the inbox without completing its ticket, checking
|
||||
* {@code fleet_reply} straight to the async ticket it belongs to whenever exactly one is still
|
||||
* parked waiting for it (see {@link #askAnsweredAsyncTasks}), so by the time a session is
|
||||
* released its tasks are normally already resolved — this loop's {@code complete} calls are then
|
||||
* harmless no-ops (a {@link CompletableFuture} can only resolve once). But should some other path
|
||||
* someday strand a reply in the inbox without completing its ticket, checking
|
||||
* {@link #hasStrandedReply(String)} here — before ever writing a failure — means a torn-down
|
||||
* session whose worker in fact replied is still reported {@code REPLIED} with that reply's own
|
||||
* text, never the misleading "the worker session was released before it replied" (which also
|
||||
* means the snapshot/worktree recovery hint that follows it never prints once a reply exists).
|
||||
*
|
||||
* <p><strong>At most one task gets the recovered reply — and here, unlike {@link #reply}'s
|
||||
* {@link #askAnsweredAsyncTasks}, {@code matching.size() >= 2} alone is reachable today.</strong>
|
||||
* This method's {@code matching} filter has no {@code turnId != null} requirement, so it matches
|
||||
* any plain (never-asked) open task too — and {@link #sendAsync} does not limit a target to one
|
||||
* of those: a second {@code fleet_send{wait:false}} at a target that is still busy returns its own
|
||||
* ticket immediately and simply parks its {@link #send} behind the target's session lock for up
|
||||
* to {@link #ASYNC_TIMEOUT_MS}, exactly as {@code abandonFailsEveryPendingAsyncTicketForTheReleasedTarget}
|
||||
* already proves. Before this fix, the loop below drained the strand once and then reused that
|
||||
* same {@code Reply} for <em>every</em> task it walked past — so two open tasks really did both
|
||||
* complete {@code REPLIED} with the same text (see the pre-fix loop in commit 97f6c33's parent).
|
||||
* A stranded reply is one worker answer, so it can settle at most one open task on this target —
|
||||
* never every open task, and never a guess. When more than one task is still open here, the
|
||||
* recovered reply goes to the <em>oldest</em> (lowest {@link Task#createdNanos}) — it has been
|
||||
* waiting longest, so it is the one most likely to be what the reply actually answers. Every
|
||||
* other open task keeps the ordinary {@code WORKER_FAILED} path it would take without a stranded
|
||||
* reply at all.
|
||||
*
|
||||
* <p><strong>{@code matching.size() >= 2} together with {@code hadStrandedReply} is a different
|
||||
* question, and today it is defence in depth rather than a path this codebase's public API can
|
||||
* drive.</strong> This class has exactly two sites that ever acquire a target's entry in
|
||||
* {@code sessionLocks} — {@link #send} and {@link #answer} — and both open a {@link Rendezvous}
|
||||
* waiter for that same target as the very first thing they do after acquiring the lock, then hold
|
||||
* lock and waiter together for the rest of their critical section ({@link #send} also clears
|
||||
* {@link #strandedReplies} right there, the instant it opens its waiter — before it ever enqueues
|
||||
* delivery). So "the session lock is held" and "a live waiter is open for it" are the same fact
|
||||
* throughout this class, and {@link #reply}'s fast path always resolves a currently-open waiter
|
||||
* directly rather than stranding. The two facts this method wants therefore cannot be produced
|
||||
* side by side: while the lock is held, a real reply resolves the open waiter directly and never
|
||||
* reaches {@link #strandedReplies}; the instant the lock is free, any parked matching task's own
|
||||
* {@link #send} that is scheduled next wins it and, by opening its waiter, clears the strand again
|
||||
* before this method ever runs. There is no way to hold that lock open-but-unaccepted from outside
|
||||
* {@link #send}/{@link #answer} to freeze a window in between. Constructing both facts at once
|
||||
* through {@code sendAsync}/{@code reply}/{@code ask}/{@code answer} would need a race against
|
||||
* virtual-thread scheduling, not a deterministic sequence — so the oldest-wins code below stays as
|
||||
* defence in depth against a regression to that mechanism (e.g. clearing {@link #strandedReplies}
|
||||
* on a narrower condition than "any acceptance"), not because today's test suite exercises the
|
||||
* conjunction. {@code matching.size() >= 2} alone, without a strand, is exactly what
|
||||
* {@code abandonFailsEveryPendingAsyncTicketForTheReleasedTarget} already covers.
|
||||
*
|
||||
* @return true if a live waiter or an async task was failed (never true for one recovered as a
|
||||
* reply — see the note above)
|
||||
*/
|
||||
@@ -494,21 +579,39 @@ public final class MessageService {
|
||||
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
|
||||
boolean failed = waiter != null && !waiter.isDone() && rendezvous.resolveFailure(waiter, reason);
|
||||
boolean asyncFailed = false;
|
||||
Reply recovered = null; // lazily drained at most once, only if a task actually needs it
|
||||
|
||||
List<Task> matching = new ArrayList<>();
|
||||
for (Task task : tasks.values()) {
|
||||
if (!target.equals(task.target) || task.question != null || task.future.isDone()) {
|
||||
continue;
|
||||
if (target.equals(task.target) && task.question == null && !task.future.isDone()) {
|
||||
matching.add(task);
|
||||
}
|
||||
if (hadStrandedReply && recovered == null) {
|
||||
recovered = recoverStrandedReply(target);
|
||||
}
|
||||
Task recoveryTask = null;
|
||||
if (hadStrandedReply && !matching.isEmpty()) {
|
||||
recoveryTask = matching.get(0);
|
||||
for (Task candidate : matching) {
|
||||
if (candidate.createdNanos < recoveryTask.createdNanos) {
|
||||
recoveryTask = candidate;
|
||||
}
|
||||
}
|
||||
Reply outcome = recovered != null ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
|
||||
}
|
||||
Reply recovered = recoveryTask != null ? recoverStrandedReply(target) : null;
|
||||
for (Task task : matching) {
|
||||
boolean isRecovery = task == recoveryTask && recovered != null;
|
||||
Reply outcome = isRecovery ? recovered : new Reply(Outcome.WORKER_FAILED, reason);
|
||||
if (task.future.complete(outcome)) {
|
||||
if (outcome.outcome() == Outcome.WORKER_FAILED) {
|
||||
asyncFailed = true;
|
||||
} else if (task.turnId != null) {
|
||||
asyncTasksByTurn.remove(task.turnId, task);
|
||||
}
|
||||
} 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());
|
||||
}
|
||||
}
|
||||
if (failed) {
|
||||
@@ -523,6 +626,10 @@ public final class MessageService {
|
||||
* stranding fact raced away, e.g. a lead's own {@code fleet_poll} on the raw session already
|
||||
* drained it first). When more than one message is queued, only the newest is the worker's actual
|
||||
* final answer ({@link #drainReplies} returns them oldest-first).
|
||||
*
|
||||
* <p>This does drain (removes the messages from the inbox) before the caller knows whether the
|
||||
* task it is recovering for will actually accept them — {@link #abandon} is the one that puts a
|
||||
* reply back if its {@code complete} call turns out to lose the race.
|
||||
*/
|
||||
private Reply recoverStrandedReply(String target) {
|
||||
var messages = drainReplies(target);
|
||||
|
||||
@@ -315,9 +315,7 @@ public final class FleetApp {
|
||||
.map(Agent.class::cast)
|
||||
.filter(a -> a.terminalId() != null)
|
||||
.collect(Collectors.toMap(Agent::terminalId, Function.identity(), (_, b) -> b));
|
||||
// fleetd #209: this REST roster reports agentSessionId via SessionManager.rosterView, so it
|
||||
// uses the resolving roster read (caller-driven, not a timer) rather than the plain one.
|
||||
List<Map<String, Object>> out = sessions.rosterResolved().stream()
|
||||
List<Map<String, Object>> out = sessions.roster().stream()
|
||||
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
|
||||
.toList();
|
||||
Map<String, Object> body = new LinkedHashMap<>();
|
||||
|
||||
@@ -89,15 +89,4 @@ public record MemberSession(
|
||||
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
nowNanos, turnCount + 1, state, worktree, branch, charterReceipt, agentSessionId);
|
||||
}
|
||||
|
||||
/**
|
||||
* Return a copy with {@code agentSessionId} resolved to a non-null value (fleetd #209). Some
|
||||
* adapters (opencode) cannot answer {@link dev.ltms.fleet.peer.PeerHandle#agentSessionId()} at
|
||||
* spawn time — the peer has not persisted its session record yet — so the id is discovered on
|
||||
* a later poll and swapped into the otherwise-immutable session via this wither.
|
||||
*/
|
||||
public MemberSession withAgentSessionId(String agentSessionId) {
|
||||
return new MemberSession(paneId, terminalId, profile, role, cwd, ownerTerminal, spawnedAtNanos,
|
||||
lastActivityAtNanos, turnCount, state, worktree, branch, charterReceipt, agentSessionId);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -46,15 +46,6 @@ public final class SessionManager implements TurnListener {
|
||||
private final PeerLauncher launcher;
|
||||
private final Worktrees worktrees;
|
||||
private final ConcurrentHashMap<String /*paneId*/, MemberSession> registry = new ConcurrentHashMap<>();
|
||||
/**
|
||||
* fleetd #209: the live {@link PeerHandle} for every registered pane, retained solely so
|
||||
* {@link #resolveAgentSessionId} can re-poll {@link PeerHandle#agentSessionId()} after spawn.
|
||||
* The handle used to go out of scope at the end of the spawn method, so a launcher that answers
|
||||
* the id lazily (opencode — the on-disk session row is written after the pane is created) could
|
||||
* never be re-asked, and {@code fleet_list}/{@code fleet_spawn resumeSessionId} never saw it.
|
||||
* Populated on every spawn path, removed on {@link #release}.
|
||||
*/
|
||||
private final ConcurrentHashMap<String /*paneId*/, PeerHandle> handles = new ConcurrentHashMap<>();
|
||||
private final MemberPresence presence;
|
||||
private final SecureRandom nonceRandom = new SecureRandom();
|
||||
private final AtomicLong nonceSeq = new AtomicLong();
|
||||
@@ -214,7 +205,6 @@ public final class SessionManager implements TurnListener {
|
||||
handle.charterReceipt(),
|
||||
handle.agentSessionId());
|
||||
registry.put(handle.id(), session);
|
||||
handles.put(handle.id(), handle);
|
||||
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
|
||||
log.debug("acquired session id={} terminal={} profile={} owner={}",
|
||||
handle.id(), handle.terminalId(), session.profile(), session.ownerTerminal());
|
||||
@@ -267,10 +257,6 @@ public final class SessionManager implements TurnListener {
|
||||
*/
|
||||
private void release(String paneId, ReleaseCause cause) {
|
||||
MemberSession removed = registry.remove(paneId);
|
||||
// fleetd #209: remove right alongside the registry entry so a released session's handle is
|
||||
// never leaked — but keep the local reference below, so the id can still be resolved for
|
||||
// the ReleaseDetail this teardown notifies with.
|
||||
PeerHandle removedHandle = handles.remove(paneId);
|
||||
boolean preserveWorktree = cause == ReleaseCause.SHUTDOWN;
|
||||
String snapshotRef = null;
|
||||
if (removed != null) {
|
||||
@@ -317,11 +303,8 @@ public final class SessionManager implements TurnListener {
|
||||
// too, so a failed ticket's detail can point a lead at the same tree to re-dispatch.
|
||||
// CB-584 (issue #65 criterion 5): carry agentSessionId alongside them, so a lead can
|
||||
// also resume the member's conversation, not just re-dispatch onto its files.
|
||||
// fleetd #209: a late-resolving adapter (opencode) may only now have an id — resolve
|
||||
// one last time so a released member's detail carries the id it now has.
|
||||
MemberSession resolved = resolveAgentSessionId(removed, removedHandle);
|
||||
notifyReleased(new ReleaseDetail(resolved.terminalId(), resolved.worktree(),
|
||||
resolved.branch(), snapshotRef, resolved.agentSessionId()));
|
||||
notifyReleased(new ReleaseDetail(removed.terminalId(), removed.worktree(),
|
||||
removed.branch(), snapshotRef, removed.agentSessionId()));
|
||||
}
|
||||
}
|
||||
// CB-581: the pane must always stop, even if the dirty check above threw. A session removed
|
||||
@@ -525,7 +508,6 @@ public final class SessionManager implements TurnListener {
|
||||
handle.charterReceipt(),
|
||||
handle.agentSessionId());
|
||||
registry.put(handle.id(), session);
|
||||
handles.put(handle.id(), handle);
|
||||
memberLifecycle.acquired(session.role(), session.profile(), session.terminalId());
|
||||
log.debug("acquired worktree session id={} terminal={} profile={} branch={} path={}",
|
||||
handle.id(), handle.terminalId(), session.profile(), session.branch(), session.worktree());
|
||||
@@ -557,73 +539,16 @@ public final class SessionManager implements TurnListener {
|
||||
return launcher.defaultProfile();
|
||||
}
|
||||
|
||||
/**
|
||||
* The session for {@code paneId}, if it is still registered and not released. fleetd #209:
|
||||
* resolves a still-unknown {@code agentSessionId} against the retained handle before returning,
|
||||
* so {@code fleet_status} sees an id a lazy-resolving adapter has since written.
|
||||
*/
|
||||
/** The session for {@code paneId}, if it is still registered and not released. */
|
||||
public Optional<MemberSession> get(String paneId) {
|
||||
return Optional.ofNullable(registry.get(paneId)).map(this::resolveAgentSessionId);
|
||||
return Optional.ofNullable(registry.get(paneId));
|
||||
}
|
||||
|
||||
/**
|
||||
* Fleet-owned roster: all registered sessions (acquired minus released). Deliberately does
|
||||
* <strong>not</strong> resolve {@code agentSessionId} (fleetd #209 follow-up) — this is the
|
||||
* roster supplier on the heartbeat and health-tick timers ({@code LeadHeartbeatLoop},
|
||||
* {@code FleetHealthMonitor} in {@code Fleetd}), on the placement/exhaustion paths, and on the
|
||||
* metrics scrape ({@code FleetMetrics}), all called far more often than any caller actually
|
||||
* reads {@code agentSessionId}. Resolving here would mean every tick opens a lazy-resolving
|
||||
* adapter's (opencode's) on-disk session store once per member whose id is still unknown — and
|
||||
* for a member whose id never appears, that cost never stops, for the life of the process. Use
|
||||
* {@link #rosterResolved()} instead wherever the id must be current.
|
||||
*/
|
||||
/** Fleet-owned roster: all registered sessions (acquired minus released). */
|
||||
public List<MemberSession> roster() {
|
||||
return List.copyOf(registry.values());
|
||||
}
|
||||
|
||||
/**
|
||||
* {@link #roster()}, with each session's still-unknown {@code agentSessionId} re-resolved
|
||||
* against its retained handle (fleetd #209) — so a caller that actually reports the id (
|
||||
* {@code fleet_list}, the REST roster) sees one a lazy-resolving adapter (opencode) has since
|
||||
* written, rather than the null frozen in at spawn time. Reserved for caller-driven reads, not
|
||||
* timers: see {@link #roster()}'s javadoc for why the plain roster must stay non-resolving.
|
||||
*/
|
||||
public List<MemberSession> rosterResolved() {
|
||||
return registry.values().stream().map(this::resolveAgentSessionId).toList();
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve {@code session}'s {@code agentSessionId} if still unknown, re-polling the retained
|
||||
* {@link PeerHandle} for this pane (fleetd #209). A no-op — returning {@code session} unchanged
|
||||
* — once the id is already known, once no handle is retained for this pane (never spawned, or
|
||||
* already released), or if the handle throws while answering. A resolved id is best-effort
|
||||
* CAS-swapped into the registry via {@link #replace}; a lost race just means another caller
|
||||
* already applied the same update, so the resolved value is returned either way.
|
||||
*/
|
||||
private MemberSession resolveAgentSessionId(MemberSession session) {
|
||||
return resolveAgentSessionId(session, handles.get(session.paneId()));
|
||||
}
|
||||
|
||||
private MemberSession resolveAgentSessionId(MemberSession session, PeerHandle handle) {
|
||||
if (session.agentSessionId() != null || handle == null) {
|
||||
return session;
|
||||
}
|
||||
String resolved;
|
||||
try {
|
||||
resolved = handle.agentSessionId();
|
||||
} catch (RuntimeException e) {
|
||||
log.debug("agentSessionId lookup failed for pane={} terminal={}: {}",
|
||||
session.paneId(), session.terminalId(), e.toString());
|
||||
return session;
|
||||
}
|
||||
if (resolved == null) {
|
||||
return session;
|
||||
}
|
||||
MemberSession updated = session.withAgentSessionId(resolved);
|
||||
replace(session, updated); // best-effort; a lost CAS just means the resolved value stands anyway
|
||||
return updated;
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-304 merged roster+live view. The registry is authoritative for worktree, branch,
|
||||
* profile, owner, and state; the optional live agent supplies the herdr-reported status.
|
||||
|
||||
@@ -688,6 +688,32 @@ class MessageServiceTest {
|
||||
assertFailedTicket(third, "agent target term_a not found");
|
||||
}
|
||||
|
||||
// --- #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
|
||||
// open matching tasks was attempted here and removed after investigation showed the combination
|
||||
// is not reachable through the public API today, not merely hard to time right:
|
||||
//
|
||||
// This class has exactly two call sites that ever hold a target's entry in the session-lock map
|
||||
// (send() and answer()), and both open a Rendezvous waiter for that same target as the first thing
|
||||
// they do after acquiring the lock, holding lock and waiter together for their whole critical
|
||||
// section. So "the lock is held" and "a live waiter is open" are the same fact throughout this
|
||||
// class. reply()'s fast path always resolves a currently-open waiter directly instead of
|
||||
// stranding — so a strand can only be created while NO task is accepted (lock free), and the
|
||||
// instant the lock is next taken (by any parked matching task's own send(), the moment it is
|
||||
// scheduled), that acceptance clears strandedReplies again (see send()'s CB-640 comment) before
|
||||
// abandon() can ever observe both facts together. Confirmed empirically too: an earlier version of
|
||||
// this test stranded a reply, then created an "accepted" task (awaitWaiting()) followed by a
|
||||
// "parked" one — and the accepted task's own acceptance silently cleared the strand it was
|
||||
// supposed to be racing against, so the parked task came back WORKER_FAILED instead of DONE, not
|
||||
// because the fix was missing but because the test's premise could not be constructed.
|
||||
//
|
||||
// The reachable half — matching.size() >= 2 alone, no strand — is exactly what
|
||||
// abandonFailsEveryPendingAsyncTicketForTheReleasedTarget already covers (all fail, none guess).
|
||||
// The oldest-wins code in abandon() stays as defence in depth (see its own javadoc) against a
|
||||
// regression that would make the conjunction reachable, e.g. clearing strandedReplies on a
|
||||
// narrower condition than "any acceptance" — not because this suite exercises it today.
|
||||
|
||||
@Test
|
||||
void abandonDoesNotFailAnAsyncTicketWaitingForAnAnswer() throws Exception {
|
||||
String ticket = messages.sendAsync(T, "task that asks");
|
||||
|
||||
@@ -931,235 +931,4 @@ class SessionManagerTest {
|
||||
assertTrue(e.getMessage().contains("SESSION_RESUME"), e.getMessage());
|
||||
assertTrue(e.getMessage().contains("stub-profile"), e.getMessage());
|
||||
}
|
||||
|
||||
// ── fleetd #209: agentSessionId resolved lazily against the retained handle ─────────────────
|
||||
//
|
||||
// The opencode adapter cannot answer PeerHandle.agentSessionId() at spawn time — the on-disk
|
||||
// session row is written only after the pane is live — so the id must be re-polled on a LATER
|
||||
// call, against the SAME handle instance the launcher returned at spawn. SessionManager used
|
||||
// to let that handle go out of scope at the end of the spawn method, so no caller ever re-asked
|
||||
// it and fleet_list/fleet_status never saw the id. LazyIdHandle below reproduces exactly that
|
||||
// shape: null on the first N calls (the spawn-time call included), a real id after.
|
||||
|
||||
/**
|
||||
* A {@link PeerHandle} whose {@link #agentSessionId()} answers {@code null} for its first
|
||||
* {@code nullCalls} invocations, then either a fixed id or a configured throw on every call
|
||||
* after that — the shape of the opencode bug (fleetd #209): the session row is not written
|
||||
* until after the pane is live, so early polls come back empty and a later one finds it.
|
||||
*/
|
||||
private static final class LazyIdHandle implements PeerHandle {
|
||||
private final String id;
|
||||
private final String terminalId;
|
||||
private final int nullCalls;
|
||||
private final String resolvedId;
|
||||
private final java.util.concurrent.atomic.AtomicInteger calls =
|
||||
new java.util.concurrent.atomic.AtomicInteger();
|
||||
private volatile RuntimeException throwAfter;
|
||||
|
||||
LazyIdHandle(String id, String terminalId, int nullCalls, String resolvedId) {
|
||||
this.id = id;
|
||||
this.terminalId = terminalId;
|
||||
this.nullCalls = nullCalls;
|
||||
this.resolvedId = resolvedId;
|
||||
}
|
||||
|
||||
/** After the null calls are exhausted, throw instead of answering the resolved id. */
|
||||
LazyIdHandle throwing(RuntimeException e) {
|
||||
this.throwAfter = e;
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return id;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String terminalId() {
|
||||
return terminalId;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String agentSessionId() {
|
||||
int n = calls.incrementAndGet();
|
||||
if (n <= nullCalls) {
|
||||
return null;
|
||||
}
|
||||
if (throwAfter != null) {
|
||||
throw throwAfter;
|
||||
}
|
||||
return resolvedId;
|
||||
}
|
||||
|
||||
@Override
|
||||
public CharterReceipt charterReceipt() {
|
||||
return null;
|
||||
}
|
||||
|
||||
int callCount() {
|
||||
return calls.get();
|
||||
}
|
||||
}
|
||||
|
||||
/** A minimal {@link PeerLauncher} that hands out pre-built {@link LazyIdHandle}s, one per spawn. */
|
||||
private static final class LazyIdLauncher implements PeerLauncher {
|
||||
private final java.util.Deque<LazyIdHandle> queued = new java.util.ArrayDeque<>();
|
||||
|
||||
LazyIdLauncher queue(LazyIdHandle handle) {
|
||||
queued.add(handle);
|
||||
return this;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilities() {
|
||||
return Set.of(Capability.WORKTREE, Capability.SESSION_RESUME);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilitiesFor(String profileName) {
|
||||
return capabilities();
|
||||
}
|
||||
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
LazyIdHandle handle = queued.poll();
|
||||
if (handle == null) {
|
||||
throw new IllegalStateException("no queued LazyIdHandle for this spawn");
|
||||
}
|
||||
return handle;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Set<String> profiles() {
|
||||
return Set.of("lazy");
|
||||
}
|
||||
|
||||
@Override
|
||||
public String defaultProfile() {
|
||||
return "lazy";
|
||||
}
|
||||
|
||||
@Override
|
||||
public String effectiveCwd(SpawnRequest req) {
|
||||
return "/cwd";
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<String> parityOverlay(String profileName) {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public List<?> list() {
|
||||
return List.of();
|
||||
}
|
||||
|
||||
@Override
|
||||
public int reapOrphanWorkers() {
|
||||
return 0;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void stop(String id) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean clearContext(String id) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void plainRosterDoesNotResolveAgentSessionId() {
|
||||
// fleetd #209 follow-up: roster() sits on the heartbeat/health-tick timers (and the metrics
|
||||
// scrape), so it must never trigger the resolve lookup — for opencode that lookup opens an
|
||||
// on-disk session database, and a member whose id never appears would pay that cost forever.
|
||||
// rosterResolved() is the one to use when a caller actually reports the id.
|
||||
LazyIdHandle handle = new LazyIdHandle("p0", "t0", 1, "oc-session-0");
|
||||
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
|
||||
MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null);
|
||||
assertEquals(1, handle.callCount(), "sanity: only the spawn-time call happened so far");
|
||||
|
||||
List<MemberSession> roster = sessions.roster();
|
||||
|
||||
assertEquals(1, roster.size());
|
||||
assertEquals(acquired.paneId(), roster.getFirst().paneId());
|
||||
assertNull(roster.getFirst().agentSessionId(), "the plain roster must not resolve the id");
|
||||
assertEquals(1, handle.callCount(),
|
||||
"roster() must never call agentSessionId() again — it sits on the heartbeat/health timers");
|
||||
}
|
||||
|
||||
@Test
|
||||
void rosterResolvedResolvesALateAgentSessionIdFromTheRetainedHandle() {
|
||||
LazyIdHandle handle = new LazyIdHandle("p1", "t1", 1, "oc-session-1");
|
||||
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
|
||||
|
||||
MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null);
|
||||
assertNull(acquired.agentSessionId(),
|
||||
"opencode has not written its session row yet at spawn time");
|
||||
|
||||
List<MemberSession> roster = sessions.rosterResolved();
|
||||
assertEquals(1, roster.size());
|
||||
assertEquals("oc-session-1", roster.getFirst().agentSessionId(),
|
||||
"fleet_list must see the id once the adapter can answer it");
|
||||
|
||||
Map<String, Object> view = SessionManager.rosterView(roster.getFirst(), null);
|
||||
assertEquals("oc-session-1", view.get("agentSessionId"),
|
||||
"rosterView renders whatever rosterResolved() resolved");
|
||||
}
|
||||
|
||||
@Test
|
||||
void getResolvesALateAgentSessionIdFromTheRetainedHandle() {
|
||||
LazyIdHandle handle = new LazyIdHandle("p2", "t2", 1, "oc-session-2");
|
||||
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
|
||||
MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null);
|
||||
|
||||
MemberSession resolved = sessions.get(acquired.paneId()).orElseThrow();
|
||||
|
||||
assertEquals("oc-session-2", resolved.agentSessionId(),
|
||||
"fleet_status (single-session lookup) must also see the late-resolved id");
|
||||
}
|
||||
|
||||
@Test
|
||||
void releaseCarriesALateResolvedAgentSessionIdIntoTheReleaseDetail() {
|
||||
LazyIdHandle handle = new LazyIdHandle("p3", "t3", 1, "oc-session-3");
|
||||
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
|
||||
MemberSession acquired = sessions.acquire("lazy", "/cwd", "/caller", null);
|
||||
java.util.List<String> released = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||
sessions.onRelease(detail -> released.add(detail.agentSessionId()));
|
||||
|
||||
sessions.release(acquired.paneId());
|
||||
|
||||
assertEquals(java.util.List.of("oc-session-3"), released,
|
||||
"a released member's detail carries the id it has since resolved, not the null "
|
||||
+ "frozen in at spawn time");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aThrowingHandleDoesNotBreakRosterResolved() {
|
||||
LazyIdHandle handle = new LazyIdHandle("p4", "t4", 1, "oc-session-4")
|
||||
.throwing(new RuntimeException("sqlite locked"));
|
||||
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
|
||||
sessions.acquire("lazy", "/cwd", "/caller", null);
|
||||
|
||||
List<MemberSession> roster = assertDoesNotThrow(sessions::rosterResolved,
|
||||
"a handle that throws resolving its id must not break the roster read");
|
||||
|
||||
assertEquals(1, roster.size());
|
||||
assertNull(roster.getFirst().agentSessionId(), "the id stays unresolved when the lookup throws");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aResolvedAgentSessionIdIsNotLookedUpAgain() {
|
||||
LazyIdHandle handle = new LazyIdHandle("p5", "t5", 1, "oc-session-5");
|
||||
SessionManager sessions = new SessionManager(new LazyIdLauncher().queue(handle));
|
||||
sessions.acquire("lazy", "/cwd", "/caller", null);
|
||||
assertEquals(1, handle.callCount(), "sanity: only the spawn-time call happened so far");
|
||||
|
||||
sessions.rosterResolved();
|
||||
assertEquals(2, handle.callCount(), "the first rosterResolved() read resolves the id");
|
||||
|
||||
sessions.rosterResolved();
|
||||
assertEquals(2, handle.callCount(), "once resolved, the id must not be looked up again");
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user