diff --git a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java index 6e7ff5e..d1e8d77 100644 --- a/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java +++ b/fleetd/src/main/java/dev/ltms/fleet/msg/MessageService.java @@ -168,14 +168,23 @@ public final class MessageService { private final String ticket; private final String target; private final CompletableFuture future = new CompletableFuture<>(); - private final long createdNanos; + /** + * When {@link #future} resolved, or {@code null} while it is still pending — the clock + * {@link #pruneTerminalTickets} measures the TTL from (#197). Deliberately a boxed + * {@code Long} rather than a {@code long} with a sentinel: {@link System#nanoTime} may + * legitimately return any value, zero and negatives included, so no numeric sentinel can mean + * "not stamped yet". Stamped by a {@code whenComplete} hook registered in the constructor, so + * every completion path stamps it — a reply, the completion fallback, a timeout, a failure, + * or an abandon on teardown — without each of those having to remember to. + */ + private volatile Long completedNanos; private volatile Reply question; private volatile String turnId; - private Task(String ticket, String target, long createdNanos) { + private Task(String ticket, String target, LongSupplier nowNanos) { this.ticket = ticket; this.target = target; - this.createdNanos = createdNanos; + future.whenComplete((reply, ex) -> completedNanos = nowNanos.getAsLong()); } } @@ -721,7 +730,7 @@ public final class MessageService { */ public String sendAsync(String target, String content, Runnable onAccepted) { String ticket = "task-" + ticketSeq.incrementAndGet(); - Task task = new Task(ticket, target, nowNanos.getAsLong()); + Task task = new Task(ticket, target, nowNanos); tasks.put(ticket, task); if (pushLoop != null) { // CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a @@ -821,12 +830,25 @@ public final class MessageService { * (or one the reminder cap already gave up on) is pruned here but never collected there, so it * lingers in {@code pendingTickets} forever and rides along on every later nudge to the same lead * — naming a ticket {@code fleet_poll} can no longer find (CB-588 follow-up). + * + *

The TTL runs from **completion**, not from creation (#197). It used to compare against + * {@code createdNanos}, which made the real collection window {@code TTL minus however long the + * task ran}: a delegation that took longer than the TTL had its reply destroyed on the first + * sweep after it landed, every time. That is the normal case here — real work runs well past ten + * minutes — and the reply lives only in {@code future}, so pruning it discards the worker's whole + * report with nothing to fall back on. Measuring from completion gives every ticket the same full + * window whatever its runtime, and still bounds {@code tasks}. + * + *

A task whose future is done but whose {@code completedNanos} is not stamped yet is left + * alone. That window is the few instructions between {@code complete()} and the constructor's + * {@code whenComplete} hook running; the next sweep collects it. */ private void pruneTerminalTickets() { long cutoff = nowNanos.getAsLong() - TICKET_TTL_NANOS; tasks.entrySet().removeIf(e -> { Task t = e.getValue(); - boolean expired = t.future.isDone() && t.createdNanos < cutoff; + Long completed = t.completedNanos; + boolean expired = t.future.isDone() && completed != null && completed < cutoff; if (expired && pushLoop != null) { pushLoop.ticketCollected(e.getKey()); } diff --git a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java index 60c15a1..dde5d87 100644 --- a/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java +++ b/fleetd/src/test/java/dev/ltms/fleet/msg/MessageServiceTest.java @@ -1057,6 +1057,74 @@ class MessageServiceTest { } } + /** + * #197: the ticket TTL must run from COMPLETION, not from creation. + * + *

It used to compare the cutoff against {@code createdNanos}, so the real window to collect a + * reply was {@code TTL minus however long the task ran}. A delegation that ran longer than the + * TTL was already past the cutoff the moment it finished, so the very next prune destroyed its + * reply — and the reply lives only in the task's future, so nothing could get it back. That is + * the normal case for real work here, not an edge case: three workers in one session ran well + * past ten minutes and two of their complete reports were lost this way. + * + *

The task below runs for longer than the whole TTL before it replies, which is exactly the + * shape that used to lose everything. Remove the fix and this fails: {@code poll} returns + * {@code null} because the ticket was pruned on arrival. + */ + @Test + void aTaskRunningLongerThanTheTtlStillKeepsItsReport() throws Exception { + java.util.concurrent.atomic.AtomicLong clock = new java.util.concurrent.atomic.AtomicLong(1_000_000_000L); + try (var wiring = wireWithPushLoop(1, 50, clock::get)) { + String slow = wiring.service().sendAsync(T, "a task that takes longer than the TTL"); + awaitWaiting(); + injectDelivery(); + + // The worker is still working, and has been for longer than the entire TTL. Nothing may + // be pruned yet — the ticket has not finished, so there is no report to keep or lose. + clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(30)); + + // Only now does it reply. Under the old clock this reply was born already expired. + assertTrue(rendezvous.resolve(T, "the long report")); + awaitTicketPhaseOn(wiring.service(), slow, MessageService.Phase.DONE); + + // A second delegation runs pruneTerminalTickets before it returns. + wiring.service().sendAsync(T, "an unrelated second task"); + + MessageService.TaskView view = wiring.service().poll(slow); + assertNotNull(view, "a ticket that completed just now must survive the prune, however " + + "long its task ran — the TTL is the window to COLLECT the report, not the " + + "budget for producing it"); + assertEquals(MessageService.Phase.DONE, view.phase()); + assertEquals("the long report", view.reply(), + "the worker's actual report must still be there, not just the ticket"); + } + } + + /** + * The other half of #197: the TTL must still bound {@code tasks}. Measuring from completion + * would be a leak if a finished ticket were then kept forever, so this pins the eviction that + * still has to happen — the same ticket, left uncollected for longer than the TTL AFTER it + * finished, is gone. + */ + @Test + void aFinishedTicketIsStillPrunedOnceTheTtlPassesSinceItFinished() throws Exception { + java.util.concurrent.atomic.AtomicLong clock = new java.util.concurrent.atomic.AtomicLong(1_000_000_000L); + try (var wiring = wireWithPushLoop(1, 50, clock::get)) { + String done = wiring.service().sendAsync(T, "a quick task"); + awaitWaiting(); + injectDelivery(); + assertTrue(rendezvous.resolve(T, "quick result")); + awaitTicketPhaseOn(wiring.service(), done, MessageService.Phase.DONE); + + // Nobody collected it, and the TTL has now passed since it FINISHED. + clock.addAndGet(MessageService.TICKET_TTL_NANOS + TimeUnit.SECONDS.toNanos(1)); + wiring.service().sendAsync(T, "an unrelated second task"); + + assertNull(wiring.service().poll(done), + "the TTL must still evict an uncollected finished ticket, or tasks grows forever"); + } + } + private MessageService.TaskView awaitTicketPhaseOn(MessageService svc, String ticket, MessageService.Phase phase) throws Exception { long deadline = System.currentTimeMillis() + 3000;