|
|
|
@@ -19,30 +19,33 @@ import java.util.stream.Collectors;
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* A status-gated push loop that nudges a lead's own herdr pane when it has uncollected work
|
|
|
|
|
* waiting: a worker reply queued with no live {@code bridge_send} to resolve it (CB-307), or an
|
|
|
|
|
* waiting: a worker reply queued with no live {@code bridge_send} to resolve it (CB-307), an
|
|
|
|
|
* async delegation ticket ({@code bridge_send(wait:false)}) that reached a terminal phase
|
|
|
|
|
* (CB-588).
|
|
|
|
|
* (CB-588), or an async ticket's worker pausing mid-turn in {@code bridge_ask} to await an answer
|
|
|
|
|
* (CB-582).
|
|
|
|
|
*
|
|
|
|
|
* <p><strong>CB-590: one schedule per lead.</strong> Both kinds of work are triggered through
|
|
|
|
|
* their own entry point — {@link #onReplyQueued(String)} and
|
|
|
|
|
* {@link #onTicketTerminal(String, String, boolean)} — but both resolve the lead that should be
|
|
|
|
|
* nudged and coalesce onto a single per-lead reminder schedule, tracked in {@link #activeLeads}.
|
|
|
|
|
* Earlier this was two independent schedules (one keyed by worker target for replies, one keyed
|
|
|
|
|
* by lead for tickets) that could both decide to inject into the same pane in the same window —
|
|
|
|
|
* a race, not routine behaviour, but the expensive kind: it interrupts the lead's live turn
|
|
|
|
|
* twice. Collapsing to one schedule per lead makes that structurally impossible: at most one
|
|
|
|
|
* scheduled tick chain is ever live for a given lead (guarded by {@link #activeLeads}'
|
|
|
|
|
* <p><strong>CB-590: one schedule per lead.</strong> All three kinds of work are triggered
|
|
|
|
|
* through their own entry point — {@link #onReplyQueued(String)},
|
|
|
|
|
* {@link #onTicketTerminal(String, String, boolean)}, and
|
|
|
|
|
* {@link #onQuestionOpened(String, String, String, String)} — but each resolves the lead that
|
|
|
|
|
* should be nudged and coalesces onto a single per-lead reminder schedule, tracked in
|
|
|
|
|
* {@link #activeLeads}. Earlier this was two independent schedules (one keyed by worker target
|
|
|
|
|
* for replies, one keyed by lead for tickets) that could both decide to inject into the same pane
|
|
|
|
|
* in the same window — a race, not routine behaviour, but the expensive kind: it interrupts the
|
|
|
|
|
* lead's live turn twice. Collapsing to one schedule per lead makes that structurally impossible:
|
|
|
|
|
* at most one scheduled tick chain is ever live for a given lead (guarded by {@link #activeLeads}'
|
|
|
|
|
* compare-and-set), so at most one {@code agents.send} to that lead's pane is ever in flight.
|
|
|
|
|
*
|
|
|
|
|
* <p>Each tick examines <em>everything</em> pending for that lead — reply targets whose inbox
|
|
|
|
|
* still holds an unacked message ({@link #pendingReplies}) and tickets not yet collected
|
|
|
|
|
* ({@link #pendingTickets}) — and sends at most one combined nudge per tick
|
|
|
|
|
* ({@link #injectNudge(String, int, int)}). Work that arrives while the lead is busy is never
|
|
|
|
|
* lost: it is re-read fresh on every tick until the lead is injectable or its own reminder cap
|
|
|
|
|
* ({@link #maxReminders}) is reached — reply and ticket work each spend from their own budget, so
|
|
|
|
|
* one source exhausting its cap does not stop nudges about the other (post-CB-590 regression fix;
|
|
|
|
|
* see {@link #decide}) — whichever the durable inbox / pending-ticket set doesn't already answer
|
|
|
|
|
* via {@code STOP}.
|
|
|
|
|
* still holds an unacked message ({@link #pendingReplies}), tickets not yet collected
|
|
|
|
|
* ({@link #pendingTickets}), and open questions not yet answered or lapsed
|
|
|
|
|
* ({@link #pendingQuestions}) — and sends at most one combined nudge per tick
|
|
|
|
|
* ({@link #injectNudge(String, int, int, int)}). Work that arrives while the lead is busy is
|
|
|
|
|
* never lost: it is re-read fresh on every tick until the lead is injectable or its own reminder
|
|
|
|
|
* cap ({@link #maxReminders}) is reached — each source spends from its own budget, so one source
|
|
|
|
|
* exhausting its cap does not stop nudges about the others (post-CB-590 regression fix; see
|
|
|
|
|
* {@link #decide}) — whichever the durable inbox / pending set doesn't already answer via
|
|
|
|
|
* {@code STOP}.
|
|
|
|
|
*/
|
|
|
|
|
public final class ReplyPushLoop {
|
|
|
|
|
|
|
|
|
@@ -57,6 +60,14 @@ public final class ReplyPushLoop {
|
|
|
|
|
/** Coalesced form, several uncollected tickets for the same lead. */
|
|
|
|
|
static final String TICKETS_NUDGE_FORMAT =
|
|
|
|
|
"%d tickets finished%s — run bridge_poll(ticket=...) for each to collect them: %s";
|
|
|
|
|
/** Singular form, one worker paused mid-turn in bridge_ask (CB-582) — names the answer call directly. */
|
|
|
|
|
static final String QUESTION_NUDGE_FORMAT =
|
|
|
|
|
"Worker %s asked a question (ticket %s) — answer it with bridge_send(turnId=\"%s\", "
|
|
|
|
|
+ "content=...) to resume its turn:\n%s";
|
|
|
|
|
/** Coalesced form, several open questions for the same lead. */
|
|
|
|
|
static final String QUESTIONS_NUDGE_FORMAT =
|
|
|
|
|
"%d workers are paused on a question — run bridge_poll(ticket=...) for each, then answer "
|
|
|
|
|
+ "with bridge_send(turnId=..., content=...): %s";
|
|
|
|
|
|
|
|
|
|
private final PrimaryRegistry primaryRegistry;
|
|
|
|
|
private final AgentControl agents;
|
|
|
|
@@ -75,7 +86,14 @@ public final class ReplyPushLoop {
|
|
|
|
|
private final ConcurrentHashMap<String, ReplyEntry> pendingReplies = new ConcurrentHashMap<>();
|
|
|
|
|
/** Tickets that have gone terminal but not yet been polled, keyed by ticket. */
|
|
|
|
|
private final ConcurrentHashMap<String, PendingTicket> pendingTickets = new ConcurrentHashMap<>();
|
|
|
|
|
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets). */
|
|
|
|
|
/**
|
|
|
|
|
* Open {@code bridge_ask} questions not yet answered or lapsed, keyed by {@code turnId}
|
|
|
|
|
* (CB-582). A question's own nudge count is tracked the same per-item way as
|
|
|
|
|
* {@link #pendingTickets} (CB-598): a fresh question keeps its source eligible regardless of
|
|
|
|
|
* how depleted an older, still-open question's count is.
|
|
|
|
|
*/
|
|
|
|
|
private final ConcurrentHashMap<String, PendingQuestion> pendingQuestions = new ConcurrentHashMap<>();
|
|
|
|
|
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets and/or questions). */
|
|
|
|
|
private final ConcurrentHashMap<String, Boolean> activeLeads = new ConcurrentHashMap<>();
|
|
|
|
|
|
|
|
|
|
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
|
|
|
@@ -154,6 +172,26 @@ public final class ReplyPushLoop {
|
|
|
|
|
.collect(Collectors.toUnmodifiableSet());
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* An open question awaiting the lead's answer: which ticket it belongs to, which worker asked,
|
|
|
|
|
* which lead to nudge, the question text, and how many nudges have named it so far (CB-598 —
|
|
|
|
|
* tracked per question, not per lead per source).
|
|
|
|
|
*/
|
|
|
|
|
private record PendingQuestion(String turnId, String ticket, String target, String lead,
|
|
|
|
|
String question, int nudgeCount) {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** Questions still open for {@code lead}, snapshotted fresh for one tick. */
|
|
|
|
|
private List<PendingQuestion> pendingQuestionsFor(String lead) {
|
|
|
|
|
return pendingQuestions.values().stream().filter(q -> lead.equals(q.lead())).toList();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** Question turnIds still open for {@code lead} — a plain snapshot for race comparison. */
|
|
|
|
|
private Set<String> pendingQuestionTurnIdsFor(String lead) {
|
|
|
|
|
return pendingQuestionsFor(lead).stream().map(PendingQuestion::turnId)
|
|
|
|
|
.collect(Collectors.toUnmodifiableSet());
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* The reply-source reminder count {@link #decide} should see for {@code lead} on this tick:
|
|
|
|
|
* the <em>minimum</em> nudge count among the reply targets currently pending for it (CB-598).
|
|
|
|
@@ -189,6 +227,15 @@ public final class ReplyPushLoop {
|
|
|
|
|
return min == Integer.MAX_VALUE ? 0 : min;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** As {@link #minReplyNudgeCountFor}, for the question source (CB-582). */
|
|
|
|
|
private int minQuestionNudgeCountFor(String lead) {
|
|
|
|
|
int min = Integer.MAX_VALUE;
|
|
|
|
|
for (PendingQuestion q : pendingQuestionsFor(lead)) {
|
|
|
|
|
min = Math.min(min, q.nudgeCount());
|
|
|
|
|
}
|
|
|
|
|
return min == Integer.MAX_VALUE ? 0 : min;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Pure decision function: examine everything pending for {@code lead} — reply targets and
|
|
|
|
|
* tickets alike — and return what the loop should do.
|
|
|
|
@@ -213,15 +260,29 @@ public final class ReplyPushLoop {
|
|
|
|
|
* @return the action the caller should take
|
|
|
|
|
*/
|
|
|
|
|
Action decide(String lead, int replyReminderCount, int ticketReminderCount) {
|
|
|
|
|
return decide(lead, replyReminderCount, ticketReminderCount, minQuestionNudgeCountFor(lead));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* As {@link #decide(String, int, int)}, with the question source (CB-582) folded in on the
|
|
|
|
|
* same footing as replies and tickets: its own eligibility (has open questions AND under its
|
|
|
|
|
* own {@link #maxReminders} budget) is enough on its own to {@link Action#INJECT}, exactly like
|
|
|
|
|
* the other two.
|
|
|
|
|
*
|
|
|
|
|
* @param questionReminderCount the lowest nudge count among questions open for this lead
|
|
|
|
|
*/
|
|
|
|
|
Action decide(String lead, int replyReminderCount, int ticketReminderCount, int questionReminderCount) {
|
|
|
|
|
boolean hasReplyWork = !pendingReplyTargetsFor(lead).isEmpty();
|
|
|
|
|
boolean hasTicketWork = !pendingTicketIdsFor(lead).isEmpty();
|
|
|
|
|
if (!hasReplyWork && !hasTicketWork) {
|
|
|
|
|
boolean hasQuestionWork = !pendingQuestionTurnIdsFor(lead).isEmpty();
|
|
|
|
|
if (!hasReplyWork && !hasTicketWork && !hasQuestionWork) {
|
|
|
|
|
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
|
|
|
|
|
return Action.STOP;
|
|
|
|
|
}
|
|
|
|
|
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
|
|
|
|
|
boolean ticketEligible = hasTicketWork && ticketReminderCount < maxReminders;
|
|
|
|
|
if (!replyEligible && !ticketEligible) {
|
|
|
|
|
boolean questionEligible = hasQuestionWork && questionReminderCount < maxReminders;
|
|
|
|
|
if (!replyEligible && !ticketEligible && !questionEligible) {
|
|
|
|
|
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
|
|
|
|
|
maxReminders, lead);
|
|
|
|
|
countNudge("exhausted");
|
|
|
|
@@ -300,6 +361,39 @@ public final class ReplyPushLoop {
|
|
|
|
|
pendingTickets.remove(ticket);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Called when an async ticket's worker pauses mid-turn in {@code bridge_ask} (CB-582): the
|
|
|
|
|
* question is now visible via {@code bridge_poll} (Phase.ASKING), but the reverse-rendezvous
|
|
|
|
|
* window it opened with (~55s default, see {@code BridgeMcp}/{@code BridgedApp}) is far shorter
|
|
|
|
|
* than a lead's normal minutes-long poll cadence — exactly the gap this closes. Resolves the
|
|
|
|
|
* delegating lead the same way {@link #onTicketTerminal} does and coalesces onto the same
|
|
|
|
|
* per-lead schedule (CB-590).
|
|
|
|
|
*
|
|
|
|
|
* @param ticket the async ticket the question belongs to (for {@code bridge_poll})
|
|
|
|
|
* @param target the worker session that asked
|
|
|
|
|
* @param turnId correlation id the lead answers with ({@code bridge_send turnId=...})
|
|
|
|
|
* @param question the question text
|
|
|
|
|
*/
|
|
|
|
|
public void onQuestionOpened(String ticket, String target, String turnId, String question) {
|
|
|
|
|
var lead = primaryRegistry.nudgeTargetFor(target);
|
|
|
|
|
if (lead.isEmpty()) {
|
|
|
|
|
log.debug("push: no lead is known to be waiting on {}'s question (turnId {}), skipping nudge",
|
|
|
|
|
target, turnId);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
pendingQuestions.put(turnId, new PendingQuestion(turnId, ticket, target, lead.get(), question, 0));
|
|
|
|
|
startOrCoalesce(lead.get());
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Called when a worker's {@code bridge_ask} resolves — answered or lapsed unanswered — so a
|
|
|
|
|
* scheduled tick never nudges about a question the lead already handled. A {@code turnId} that
|
|
|
|
|
* was never pending (never nudged, or already closed) is a no-op.
|
|
|
|
|
*/
|
|
|
|
|
public void questionClosed(String turnId) {
|
|
|
|
|
pendingQuestions.remove(turnId);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// --- the schedule ----------------------------------------------------------------------------
|
|
|
|
|
|
|
|
|
|
/** Start a reminder schedule for {@code lead}, or join the one already running. */
|
|
|
|
@@ -330,17 +424,19 @@ public final class ReplyPushLoop {
|
|
|
|
|
void tick(String lead) {
|
|
|
|
|
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
|
|
|
|
|
Set<String> ticketsBefore = pendingTicketIdsFor(lead);
|
|
|
|
|
Set<String> questionsBefore = pendingQuestionTurnIdsFor(lead);
|
|
|
|
|
int replyReminderCount = minReplyNudgeCountFor(lead);
|
|
|
|
|
int ticketReminderCount = minTicketNudgeCountFor(lead);
|
|
|
|
|
var action = decide(lead, replyReminderCount, ticketReminderCount);
|
|
|
|
|
int questionReminderCount = minQuestionNudgeCountFor(lead);
|
|
|
|
|
var action = decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount);
|
|
|
|
|
switch (action) {
|
|
|
|
|
case INJECT -> {
|
|
|
|
|
injectNudge(lead, replyReminderCount, ticketReminderCount);
|
|
|
|
|
injectNudge(lead, replyReminderCount, ticketReminderCount, questionReminderCount);
|
|
|
|
|
scheduleNext(lead);
|
|
|
|
|
}
|
|
|
|
|
// Re-check after the configured backoff; the lead may become injectable soon.
|
|
|
|
|
case WAIT_BUSY -> scheduleNext(lead);
|
|
|
|
|
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore);
|
|
|
|
|
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@@ -375,9 +471,20 @@ public final class ReplyPushLoop {
|
|
|
|
|
* side always wins; neither can miss the other, so this never loops on its own account.
|
|
|
|
|
*/
|
|
|
|
|
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore) {
|
|
|
|
|
stopOrRestart(lead, repliesBefore, ticketsBefore, pendingQuestionTurnIdsFor(lead));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* As {@link #stopOrRestart(String, Set, Set)}, with the question source's (CB-582) own "before"
|
|
|
|
|
* snapshot folded into the same race check: a question that raced in during the
|
|
|
|
|
* decision-to-release window reclaims the schedule slot exactly like a raced-in reply or ticket.
|
|
|
|
|
*/
|
|
|
|
|
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore,
|
|
|
|
|
Set<String> questionsBefore) {
|
|
|
|
|
activeLeads.remove(lead);
|
|
|
|
|
boolean racedIn = pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t))
|
|
|
|
|
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t));
|
|
|
|
|
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t))
|
|
|
|
|
|| pendingQuestionTurnIdsFor(lead).stream().anyMatch(t -> !questionsBefore.contains(t));
|
|
|
|
|
if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
|
|
|
|
|
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
|
|
|
|
|
scheduleNext(lead);
|
|
|
|
@@ -387,36 +494,44 @@ public final class ReplyPushLoop {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** Send one combined nudge covering everything currently pending for {@code lead}. */
|
|
|
|
|
private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount) {
|
|
|
|
|
// Re-read rather than threading it down from decide(): a reply can drain, or a ticket be
|
|
|
|
|
// collected (or another arrive), between the decision and the injection.
|
|
|
|
|
private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount,
|
|
|
|
|
int questionReminderCount) {
|
|
|
|
|
// Re-read rather than threading it down from decide(): a reply can drain, a ticket be
|
|
|
|
|
// collected, or a question be answered (or another arrive), between the decision and the
|
|
|
|
|
// injection.
|
|
|
|
|
Set<String> replyTargets = pendingReplyTargetsFor(lead);
|
|
|
|
|
List<PendingTicket> tickets = pendingTicketsFor(lead);
|
|
|
|
|
if (replyTargets.isEmpty() && tickets.isEmpty()) {
|
|
|
|
|
List<PendingQuestion> questions = pendingQuestionsFor(lead);
|
|
|
|
|
if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty()) {
|
|
|
|
|
log.debug("push: pending work for lead {} drained before the nudge could be sent", lead);
|
|
|
|
|
return;
|
|
|
|
|
}
|
|
|
|
|
String nudge = formatNudge(replyTargets, tickets);
|
|
|
|
|
String nudge = formatNudge(replyTargets, tickets, questions);
|
|
|
|
|
try {
|
|
|
|
|
agents.send(lead, nudge);
|
|
|
|
|
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}; {} reply target(s), {} ticket(s))",
|
|
|
|
|
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; "
|
|
|
|
|
+ "{} reply target(s), {} ticket(s), {} question(s))",
|
|
|
|
|
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
|
|
|
|
|
replyTargets.size(), tickets.size());
|
|
|
|
|
questionReminderCount + 1, maxReminders,
|
|
|
|
|
replyTargets.size(), tickets.size(), questions.size());
|
|
|
|
|
countNudge("delivered");
|
|
|
|
|
} catch (RuntimeException e) {
|
|
|
|
|
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}): {}",
|
|
|
|
|
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, e.toString());
|
|
|
|
|
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}, question {}/{}): {}",
|
|
|
|
|
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
|
|
|
|
|
questionReminderCount + 1, maxReminders, e.toString());
|
|
|
|
|
}
|
|
|
|
|
// Bump every item actually named in this nudge, not just whatever the shared source-level
|
|
|
|
|
// eligibility used to gate (CB-598) — each item's own count is what the next tick's
|
|
|
|
|
// minReplyNudgeCountFor / minTicketNudgeCountFor will read. An item already at or over the
|
|
|
|
|
// cap keeps riding along in the text (still pending, still named) but its extra bumps here
|
|
|
|
|
// are inert: decide() already treats it as ineligible once its count reaches maxReminders.
|
|
|
|
|
bumpNudgeCounts(replyTargets, tickets);
|
|
|
|
|
// minReplyNudgeCountFor / minTicketNudgeCountFor / minQuestionNudgeCountFor will read. An
|
|
|
|
|
// item already at or over the cap keeps riding along in the text (still pending, still
|
|
|
|
|
// named) but its extra bumps here are inert: decide() already treats it as ineligible once
|
|
|
|
|
// its count reaches maxReminders.
|
|
|
|
|
bumpNudgeCounts(replyTargets, tickets, questions);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** Record that every one of these items was just named in a sent (or attempted) nudge. */
|
|
|
|
|
private void bumpNudgeCounts(Set<String> replyTargets, List<PendingTicket> tickets) {
|
|
|
|
|
private void bumpNudgeCounts(Set<String> replyTargets, List<PendingTicket> tickets,
|
|
|
|
|
List<PendingQuestion> questions) {
|
|
|
|
|
for (String target : replyTargets) {
|
|
|
|
|
pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1));
|
|
|
|
|
}
|
|
|
|
@@ -424,6 +539,11 @@ public final class ReplyPushLoop {
|
|
|
|
|
pendingTickets.computeIfPresent(ticket.ticket(),
|
|
|
|
|
(id, e) -> new PendingTicket(e.ticket(), e.lead(), e.failed(), e.nudgeCount() + 1));
|
|
|
|
|
}
|
|
|
|
|
for (PendingQuestion question : questions) {
|
|
|
|
|
pendingQuestions.computeIfPresent(question.turnId(), (id, e) ->
|
|
|
|
|
new PendingQuestion(e.turnId(), e.ticket(), e.target(), e.lead(), e.question(),
|
|
|
|
|
e.nudgeCount() + 1));
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** Schedule the next tick on the scheduler thread pool. */
|
|
|
|
@@ -435,7 +555,8 @@ public final class ReplyPushLoop {
|
|
|
|
|
// --- nudge formatting ------------------------------------------------------------------------
|
|
|
|
|
|
|
|
|
|
/** Render everything pending for one lead as a single nudge line. */
|
|
|
|
|
private static String formatNudge(Set<String> replyTargets, List<PendingTicket> tickets) {
|
|
|
|
|
private static String formatNudge(Set<String> replyTargets, List<PendingTicket> tickets,
|
|
|
|
|
List<PendingQuestion> questions) {
|
|
|
|
|
List<String> parts = new ArrayList<>();
|
|
|
|
|
if (!replyTargets.isEmpty()) {
|
|
|
|
|
parts.add(formatRepliesNudge(replyTargets));
|
|
|
|
@@ -443,6 +564,9 @@ public final class ReplyPushLoop {
|
|
|
|
|
if (!tickets.isEmpty()) {
|
|
|
|
|
parts.add(formatTicketsNudge(tickets));
|
|
|
|
|
}
|
|
|
|
|
if (!questions.isEmpty()) {
|
|
|
|
|
parts.add(formatQuestionsNudge(questions));
|
|
|
|
|
}
|
|
|
|
|
return String.join(" | ", parts);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@@ -470,6 +594,18 @@ public final class ReplyPushLoop {
|
|
|
|
|
return TICKETS_NUDGE_FORMAT.formatted(pending.size(), failedNote, ids);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** Render one or several open questions (CB-582). */
|
|
|
|
|
private static String formatQuestionsNudge(List<PendingQuestion> pending) {
|
|
|
|
|
if (pending.size() == 1) {
|
|
|
|
|
PendingQuestion q = pending.get(0);
|
|
|
|
|
return QUESTION_NUDGE_FORMAT.formatted(q.target(), q.ticket(), q.turnId(), q.question());
|
|
|
|
|
}
|
|
|
|
|
String ids = pending.stream()
|
|
|
|
|
.map(q -> q.ticket() + " (turnId=" + q.turnId() + ")")
|
|
|
|
|
.collect(Collectors.joining(", "));
|
|
|
|
|
return QUESTIONS_NUDGE_FORMAT.formatted(pending.size(), ids);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// --- lifecycle -----------------------------------------------------------------------------
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
@@ -477,8 +613,8 @@ public final class ReplyPushLoop {
|
|
|
|
|
* uses this to stand aside: while the push loop is actively nudging a lead, a concurrent
|
|
|
|
|
* heartbeat injection would start a second competing turn in the same pane — racing loops
|
|
|
|
|
* multiply turns and context burn. "Active" means a schedule exists in {@link #activeLeads},
|
|
|
|
|
* which now covers both reply-queued (CB-307) and ticket-terminal (CB-588) work (CB-590) —
|
|
|
|
|
* bounded by what has been triggered, not by any persistent state.
|
|
|
|
|
* which now covers reply-queued (CB-307), ticket-terminal (CB-588), and question-open (CB-582)
|
|
|
|
|
* work (CB-590) — bounded by what has been triggered, not by any persistent state.
|
|
|
|
|
*/
|
|
|
|
|
public boolean isActive() {
|
|
|
|
|
return !activeLeads.isEmpty();
|
|
|
|
@@ -490,6 +626,7 @@ public final class ReplyPushLoop {
|
|
|
|
|
activeLeads.clear();
|
|
|
|
|
pendingReplies.clear();
|
|
|
|
|
pendingTickets.clear();
|
|
|
|
|
pendingQuestions.clear();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/** @see #stop() */
|
|
|
|
|