A TIMED_OUT_QUEUED send leaves its message in the injector queue, and it is delivered later as a stale brief #338

Closed
opened 2026-09-04 10:47:42 +02:00 by ltms · 1 comment
Owner

Found by a hunt over inject/. I have re-read the two call sites myself and confirmed the path.
This one is not theory — it is a behaviour I have hit while operating the fleet.

What happens

MessageService.send enqueues into the injector, then waits on the rendezvous:

// MessageService.java:866
CompletableFuture<Void> delivered = injector.enqueue(target, content, token);
try {
    Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), ...);
    ...
} catch (TimeoutException e) {
    boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally();
    if (!wasDelivered) {
        // records a health fact only
        queuedDeliveries.put(target, Boolean.TRUE);
    }
    return recorded(new Reply(
            wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null));
}
...
} finally {
    asyncTasksByWaiter.remove(reply);
    rendezvous.close(target, reply);   // closes the WAITER
}

On TIMED_OUT_QUEUED the caller is told the send did not go. The finally closes the rendezvous
waiter. Nothing removes the Pending from the injector's queue.

Later, when the member becomes injectable, Injector.onStatus picks it straight back up:

// Injector.java:272-280
Pending p = t.queue.peek();
if (p != null && ready.test(target)) {
    agentsFor(target).send(target, p.text());
    t.queue.poll();
    ...
}

So the old text is typed into the member's terminal, minutes or hours after its sender was told it
had timed out.

Direction of harm — this is the loose direction, and it is the bad one

queuedDeliveries records the fact for health reporting. It is not a cancellation.
CompletionResolver avoids resolving the already-closed waiter, so nothing wrong is reported back —
but nothing stops AgentControl.send either. The message still lands.

The lead has been told "queued, not delivered", has usually moved on, and often re-sent the work
elsewhere. Then the member wakes up and starts on the stale brief. Nobody is told.

Operational confirmation

I have hit this while running the fleet, before this hunt, and wrote it down at the time: a second
send to a busy member is reported as timed_out_queued and never delivered, and the stale brief is
still queued and restarts the member later when it goes idle. The hunt found the code behind the
symptom I already had.

MessageServiceTest.java:1764-1772 documents the queued-and-stays-queued behaviour, so this is not
an accident of the current code — it is asserted. Whether that assertion is still the behaviour we
want is part of this ticket.

Goal and invariants

Goal: when a send's deadline passes and the message was never delivered, that message must not
be delivered afterwards. The sender was told it did not go; the fleet must make that true.

Invariants:

  1. TIMED_OUT_WORKING is a different case and must not change. There the message WAS delivered and
    the member is working on it — cancelling anything there would lose real work.
  2. Do not break queue ordering for the messages that remain. Cancelling the middle of a queue must
    not strand what is behind it.
  3. queuedDeliveries' health reporting stays. It is orthogonal — it says what happened, and this
    ticket changes what happens next.

Candidate mechanism, as a candidate only: enqueue already takes a token. Give the sender a way
to cancel that exact Pending — matched by identity, not by target or by text — and call it from
the !wasDelivered branch. Decide it yourself and justify it. In particular say what happens if the
cancel races the injector picking that same Pending up: at that instant the message IS being
delivered, and the honest answer may be to let it go and report TIMED_OUT_WORKING instead. Say
which you chose.

Also update MessageServiceTest.java:1764-1772, which currently asserts the old behaviour. A
test asserting the bug is part of the bug — change it deliberately and say so in the PR body, do not
delete it quietly.

Direction of harm to weigh, honestly

Cancelling is the strict direction and it has its own cost: a member that was about to become idle
loses a message it could have taken, and the lead must re-send. That is loud and recoverable. The
current behaviour is silent and is not. Prefer loud — but say in the PR body what the strict
direction costs, rather than only what it fixes.

Found by a hunt over `inject/`. I have re-read the two call sites myself and confirmed the path. **This one is not theory — it is a behaviour I have hit while operating the fleet.** ## What happens `MessageService.send` enqueues into the injector, then waits on the rendezvous: ```java // MessageService.java:866 CompletableFuture<Void> delivered = injector.enqueue(target, content, token); try { Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), ...); ... } catch (TimeoutException e) { boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally(); if (!wasDelivered) { // records a health fact only queuedDeliveries.put(target, Boolean.TRUE); } return recorded(new Reply( wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null)); } ... } finally { asyncTasksByWaiter.remove(reply); rendezvous.close(target, reply); // closes the WAITER } ``` On `TIMED_OUT_QUEUED` the caller is told the send did not go. The `finally` closes the rendezvous waiter. **Nothing removes the `Pending` from the injector's queue.** Later, when the member becomes injectable, `Injector.onStatus` picks it straight back up: ```java // Injector.java:272-280 Pending p = t.queue.peek(); if (p != null && ready.test(target)) { agentsFor(target).send(target, p.text()); t.queue.poll(); ... } ``` So the old text is typed into the member's terminal, minutes or hours after its sender was told it had timed out. ## Direction of harm — this is the loose direction, and it is the bad one `queuedDeliveries` records the fact for health reporting. It is **not** a cancellation. `CompletionResolver` avoids resolving the already-closed waiter, so nothing wrong is reported back — but nothing stops `AgentControl.send` either. The message still lands. The lead has been told "queued, not delivered", has usually moved on, and often re-sent the work elsewhere. Then the member wakes up and starts on the stale brief. Nobody is told. ## Operational confirmation I have hit this while running the fleet, before this hunt, and wrote it down at the time: a second send to a busy member is reported as `timed_out_queued` and never delivered, and the stale brief is still queued and restarts the member later when it goes idle. The hunt found the code behind the symptom I already had. `MessageServiceTest.java:1764-1772` documents the queued-and-stays-queued behaviour, so this is not an accident of the current code — it is asserted. Whether that assertion is still the behaviour we want is part of this ticket. ## Goal and invariants **Goal:** when a send's deadline passes and the message was never delivered, that message must not be delivered afterwards. The sender was told it did not go; the fleet must make that true. **Invariants:** 1. `TIMED_OUT_WORKING` is a different case and must not change. There the message WAS delivered and the member is working on it — cancelling anything there would lose real work. 2. Do not break queue ordering for the messages that remain. Cancelling the middle of a queue must not strand what is behind it. 3. `queuedDeliveries`' health reporting stays. It is orthogonal — it says what happened, and this ticket changes what happens next. **Candidate mechanism, as a candidate only:** `enqueue` already takes a token. Give the sender a way to cancel that exact `Pending` — matched by identity, not by target or by text — and call it from the `!wasDelivered` branch. Decide it yourself and justify it. In particular say what happens if the cancel races the injector picking that same `Pending` up: at that instant the message IS being delivered, and the honest answer may be to let it go and report `TIMED_OUT_WORKING` instead. Say which you chose. **Also update `MessageServiceTest.java:1764-1772`**, which currently asserts the old behaviour. A test asserting the bug is part of the bug — change it deliberately and say so in the PR body, do not delete it quietly. ## Direction of harm to weigh, honestly Cancelling is the strict direction and it has its own cost: a member that was about to become idle loses a message it could have taken, and the lead must re-send. That is loud and recoverable. The current behaviour is silent and is not. Prefer loud — but say in the PR body what the strict direction costs, rather than only what it fixes.
Author
Owner

Merged to main as a8cadd9. The branch was up to date with main — checked with
git merge-base --is-ancestor, not assumed.

Build after the merge, unpiped: Tests run: 1357, Failures: 0, Errors: 0, Skipped: 0,
BUILD SUCCESS.

The race question was answered properly

This was the part of the ticket that could have gone wrong quietly, and it did not. The timeout
branch now reads the cancellation result:

if (!wasDelivered) {
    // The target monitor makes cancellation atomic with onStatus picking this
    // Pending up. If pickup won, report TIMED_OUT_WORKING because the text landed.
    wasDelivered = injector.cancel(delivery) == Injector.Cancellation.DELIVERED;
}

So a lost race reports TIMED_OUT_WORKING, not TIMED_OUT_QUEUED. Reporting queued while the text
landed would have been the same bug with a smaller window.

The test that asserted the old behaviour

MessageServiceTest:1768 did assert the bug. The implementer kept the health assertion, corrected
its misleading message, and added the assertion that matters:

assertTrue(messages.hasQueuedDelivery(T),
        "a TIMED_OUT_QUEUED send still records the undelivered delivery for fleet health");
injector.onStatus(T, AgentStatus.IDLE);
assertTrue(herdr.calls.stream().noneMatch(c -> c.method().equals("agent.prompt")),
        "a TIMED_OUT_QUEUED send must be cancelled, not delivered when the worker later goes idle");

That is the deliberate change the ticket asked for, and the new assertion is the bug itself.

My own mutation — aimed at a different line

The implementer proved its fix by removing the cancellation. I mutated the race resolution
instead.

Mutation W — ignore cancel's return value (injector.cancel(delivery);, leaving
wasDelivered false):

Tests run: 1357, Failures: 0, Errors: 0, Skipped: 0
BUILD SUCCESS

Nothing failed. InjectorTest:190-197 pins Injector.cancel returning DELIVERED, so the seam
is covered — but nothing pins MessageService reading it. Seam tested, caller not.

Injector is final, so a test cannot stub cancel for this, and driving the real interleaving
needs a fourth test-only hook in a class that already has three. I did not add one on my own
judgment. Filed as #345 with the three options and their trade-offs. Reverted; main is unchanged.

The cost is stated

The PR body says plainly what strict cancellation costs: a member about to go idle loses a message
it could have taken, and the lead must send again. Loud and recoverable, against silent and not.
That is the right trade, and it is written down rather than assumed.

Closing. One fix in, one coverage gap out (#345).

Merged to `main` as `a8cadd9`. The branch was up to date with main — checked with `git merge-base --is-ancestor`, not assumed. Build after the merge, unpiped: `Tests run: 1357, Failures: 0, Errors: 0, Skipped: 0`, `BUILD SUCCESS`. ## The race question was answered properly This was the part of the ticket that could have gone wrong quietly, and it did not. The timeout branch now reads the cancellation result: ```java if (!wasDelivered) { // The target monitor makes cancellation atomic with onStatus picking this // Pending up. If pickup won, report TIMED_OUT_WORKING because the text landed. wasDelivered = injector.cancel(delivery) == Injector.Cancellation.DELIVERED; } ``` So a lost race reports `TIMED_OUT_WORKING`, not `TIMED_OUT_QUEUED`. Reporting queued while the text landed would have been the same bug with a smaller window. ## The test that asserted the old behaviour `MessageServiceTest:1768` did assert the bug. The implementer kept the health assertion, corrected its misleading message, and added the assertion that matters: ```java assertTrue(messages.hasQueuedDelivery(T), "a TIMED_OUT_QUEUED send still records the undelivered delivery for fleet health"); injector.onStatus(T, AgentStatus.IDLE); assertTrue(herdr.calls.stream().noneMatch(c -> c.method().equals("agent.prompt")), "a TIMED_OUT_QUEUED send must be cancelled, not delivered when the worker later goes idle"); ``` That is the deliberate change the ticket asked for, and the new assertion is the bug itself. ## My own mutation — aimed at a different line The implementer proved its fix by removing the cancellation. I mutated the **race resolution** instead. **Mutation W — ignore `cancel`'s return value** (`injector.cancel(delivery);`, leaving `wasDelivered` false): ``` Tests run: 1357, Failures: 0, Errors: 0, Skipped: 0 BUILD SUCCESS ``` Nothing failed. `InjectorTest:190-197` pins `Injector.cancel` returning `DELIVERED`, so the **seam** is covered — but nothing pins `MessageService` reading it. Seam tested, caller not. `Injector` is `final`, so a test cannot stub `cancel` for this, and driving the real interleaving needs a fourth test-only hook in a class that already has three. I did not add one on my own judgment. Filed as #345 with the three options and their trade-offs. Reverted; main is unchanged. ## The cost is stated The PR body says plainly what strict cancellation costs: a member about to go idle loses a message it could have taken, and the lead must send again. Loud and recoverable, against silent and not. That is the right trade, and it is written down rather than assumed. Closing. One fix in, one coverage gap out (#345).
ltms closed this issue 2026-09-04 11:01:27 +02:00
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: fleet/fleetd#338