31e34d177b
A 5-turn primary↔worker conversation test (see the e2e harness) surfaced three delegation-channel gaps; this closes them. CB-115 — status + scrape correctness: - AgentStatus gains DONE (herdr's explicit turn-complete marker) so a finished turn is no longer misread as UNKNOWN and left to wedge or false-fail. - StatusRefiner reclassifies content-bearing UNKNOWN samples (StatusPoller wired to it), and the Injector baselines pane content on delivery (TurnListener gains onDelivered) to guard completion against previous-turn misattribution. - CompletionResolver.lastAssistantBlock stops at the first hard TUI boundary, so a scrape returns only the assistant answer — no input box, prompt echo, spinner, tips or warnings. CB-116 — waiter identity (the cross-turn stale reply): - The completion/failure fallback ran on a virtual thread and resolved whichever waiter was currently registered for the session. Since the rendezvous holds one waiter per session and sends serialize, turn N's late completion could land on turn N+1's waiter and deliver turn N's stale scrape as turn N+1's answer. The baseline guard missed it because turn N was resolved by bridge_reply, which never updates the completion baseline. - Fix: capture the exact waiter (and pre-turn baseline) when a turn is delivered, on the poller thread before any next-turn delivery can overwrite it, and resolve THAT waiter — a no-op if it was already resolved. Rendezvous.resolveCompletion/ resolveFailure now take the captured CompletableFuture; currentWaiter exposes the registered one for capture. A late completion for turn N can no longer touch turn N+1's send. Verified: 128 unit tests green (incl. a CB-116 regression asserting a late completion never resolves the next turn's waiter); a re-run of the conversation test passes with turn 5 resolving to its own reply rather than turn 4's text.
94 lines
3.6 KiB
Java
94 lines
3.6 KiB
Java
package dev.ltms.bridged.inject;
|
|
|
|
import dev.ltms.bridged.herdr.AgentControl;
|
|
import dev.ltms.bridged.herdr.AgentStatus;
|
|
import dev.ltms.bridged.herdr.HerdrException;
|
|
import org.slf4j.Logger;
|
|
import org.slf4j.LoggerFactory;
|
|
|
|
import java.util.Set;
|
|
|
|
/**
|
|
* Drives the {@link Injector} by sampling each active worker's {@code agent_status} and
|
|
* feeding it in. A single virtual-thread loop polls only targets that have work outstanding,
|
|
* so an idle bridge does no herdr traffic.
|
|
*
|
|
* <p>This is the Stage-1 gate signal. It can later be replaced (or fronted) by a herdr
|
|
* {@code events.subscribe} stream without touching the {@link Injector} — the injector only
|
|
* consumes {@code onStatus} calls, however they are produced.
|
|
*/
|
|
public final class StatusPoller {
|
|
|
|
private static final Logger log = LoggerFactory.getLogger(StatusPoller.class);
|
|
|
|
private final AgentControl agents;
|
|
private final Injector injector;
|
|
private final StatusRefiner refiner;
|
|
private final long intervalMillis;
|
|
private volatile boolean running;
|
|
private Thread thread;
|
|
|
|
public StatusPoller(AgentControl agents, Injector injector, long intervalMillis) {
|
|
this(agents, injector, new StatusRefiner(agents), intervalMillis);
|
|
}
|
|
|
|
public StatusPoller(AgentControl agents, Injector injector, StatusRefiner refiner,
|
|
long intervalMillis) {
|
|
this.agents = agents;
|
|
this.injector = injector;
|
|
this.refiner = refiner;
|
|
this.intervalMillis = intervalMillis;
|
|
}
|
|
|
|
/** Start the polling loop on a virtual thread. Idempotent. */
|
|
public synchronized void start() {
|
|
if (running) return;
|
|
running = true;
|
|
thread = Thread.ofVirtual().name("status-poller").start(this::loop);
|
|
log.info("status poller started (interval {}ms)", intervalMillis);
|
|
}
|
|
|
|
private void loop() {
|
|
while (running) {
|
|
Set<String> active = injector.activeTargets();
|
|
for (String target : active) {
|
|
if (!running) return;
|
|
try {
|
|
// herdr's agent_status can misreport a settled worker as `unknown`; refine it
|
|
// against the pane content before it drives delivery/completion (CB-115).
|
|
AgentStatus status = refiner.refine(target, agents.status(target));
|
|
injector.onStatus(target, status);
|
|
} catch (HerdrException e) {
|
|
// The worker's agent is gone — stop trying and unblock its waiters.
|
|
if (e.code() != null && e.code().endsWith("_not_found")) {
|
|
log.debug("target {} gone; dropping its queue", target);
|
|
injector.drop(target, e);
|
|
} else {
|
|
log.debug("status poll for {} failed (will retry): {}", target, e.getMessage());
|
|
}
|
|
} catch (RuntimeException e) {
|
|
// Never let one target's unexpected error (e.g. an odd agent.get shape) kill
|
|
// the single poller thread and stall injection for every worker.
|
|
log.warn("unexpected error polling {}; skipping this round", target, e);
|
|
}
|
|
}
|
|
sleep();
|
|
}
|
|
}
|
|
|
|
private void sleep() {
|
|
try {
|
|
Thread.sleep(intervalMillis);
|
|
} catch (InterruptedException e) {
|
|
Thread.currentThread().interrupt();
|
|
running = false;
|
|
}
|
|
}
|
|
|
|
/** Stop the polling loop. Idempotent. */
|
|
public synchronized void stop() {
|
|
running = false;
|
|
if (thread != null) thread.interrupt();
|
|
}
|
|
}
|