From e97e06f3c401920d91cf916991df533f9356aee7 Mon Sep 17 00:00:00 2001 From: "omegent-app[bot]" <306514130+omegent-app[bot]@users.noreply.github.com> Date: Mon, 27 Jul 2026 08:40:39 +0000 Subject: [PATCH] fix(discord-bot): make steernow resilient to empty/lagging server queue steernow was reporting "Nothing is queued" whenever the HTTP thread snapshot was missing or lagging, even when Discord still had parked follow-ups in the in-memory registry. Prefer the server queue, fall back to the local registry, register failed mid-turn steers for later flush, and clarify the empty-queue reply so users know to use /agent steer. --- .../DiscordQueuedPromptRegistry.test.ts | 70 ++++++++++++++++++- .../features/DiscordQueuedPromptRegistry.ts | 70 +++++++++++++++++++ .../discord-bot/src/features/MentionRouter.ts | 64 ++++++++++++++--- 3 files changed, 192 insertions(+), 12 deletions(-) diff --git a/apps/discord-bot/src/features/DiscordQueuedPromptRegistry.test.ts b/apps/discord-bot/src/features/DiscordQueuedPromptRegistry.test.ts index c68b4ffc573..985e31f7da2 100644 --- a/apps/discord-bot/src/features/DiscordQueuedPromptRegistry.test.ts +++ b/apps/discord-bot/src/features/DiscordQueuedPromptRegistry.test.ts @@ -1,7 +1,11 @@ import { MessageId, ThreadId } from "@t3tools/contracts"; import { describe, expect, it } from "vite-plus/test"; -import { createDiscordQueuedPromptRegistry } from "./DiscordQueuedPromptRegistry.ts"; +import { + createDiscordQueuedPromptRegistry, + formatSteernowEmptyQueueMessage, + resolveSteernowMessageIds, +} from "./DiscordQueuedPromptRegistry.ts"; describe("DiscordQueuedPromptRegistry", () => { it("remembers entries by Discord message and thread", () => { @@ -60,3 +64,67 @@ describe("DiscordQueuedPromptRegistry", () => { expect(registry.listForThread(ThreadId.make("thread-1"))).toEqual([]); }); }); + +describe("resolveSteernowMessageIds", () => { + it("prefers the server queue when present", () => { + const result = resolveSteernowMessageIds({ + serverQueued: [ + { messageId: MessageId.make("server-1") }, + { messageId: MessageId.make("server-2") }, + ], + localPending: [{ t3MessageId: MessageId.make("local-1") }], + detailLoaded: true, + }); + expect(result).toEqual({ + messageIds: [MessageId.make("server-1"), MessageId.make("server-2")], + source: "server", + snapshotMissing: false, + }); + }); + + it("falls back to the local registry when the server queue is empty", () => { + const result = resolveSteernowMessageIds({ + serverQueued: [], + localPending: [ + { t3MessageId: MessageId.make("local-1") }, + { t3MessageId: MessageId.make("local-1") }, + { t3MessageId: MessageId.make("local-2") }, + ], + detailLoaded: false, + }); + expect(result).toEqual({ + messageIds: [MessageId.make("local-1"), MessageId.make("local-2")], + source: "local", + snapshotMissing: true, + }); + }); + + it("reports empty with snapshotMissing when both sources are empty and detail failed", () => { + expect( + resolveSteernowMessageIds({ + serverQueued: [], + localPending: [], + detailLoaded: false, + }), + ).toEqual({ + messageIds: [], + source: "empty", + snapshotMissing: true, + }); + }); +}); + +describe("formatSteernowEmptyQueueMessage", () => { + it("mentions steer path when the queue is truly empty", () => { + const text = formatSteernowEmptyQueueMessage({ snapshotMissing: false }); + expect(text).toContain("Nothing is queued"); + expect(text).toContain("/agent steer"); + expect(text).toContain("/agent steernow"); + }); + + it("mentions snapshot failure when the HTTP detail is missing", () => { + const text = formatSteernowEmptyQueueMessage({ snapshotMissing: true }); + expect(text).toContain("thread snapshot unavailable"); + expect(text).toContain("/agent steer"); + }); +}); diff --git a/apps/discord-bot/src/features/DiscordQueuedPromptRegistry.ts b/apps/discord-bot/src/features/DiscordQueuedPromptRegistry.ts index 10132ac1f6c..af358b04ce7 100644 --- a/apps/discord-bot/src/features/DiscordQueuedPromptRegistry.ts +++ b/apps/discord-bot/src/features/DiscordQueuedPromptRegistry.ts @@ -17,6 +17,76 @@ export type PendingQueuedPrompt = { /** Discord unicode used as the "queued" badge on the user's message. */ export const QUEUED_PROMPT_REACTION_EMOJI = "📥"; +/** + * Resolve which message ids `/omegent steernow` should inject. + * + * Prefer the server queue (authoritative after restart). Fall back to the + * in-memory Discord registry when the HTTP snapshot is missing/lagging so a + * just-queued mid-turn follow-up is still steerable. + */ +export function resolveSteernowMessageIds(input: { + readonly serverQueued: ReadonlyArray<{ readonly messageId: MessageId }>; + readonly localPending: ReadonlyArray<{ readonly t3MessageId: MessageId }>; + /** True when `fetchThreadDetail` returned a snapshot (even if queue empty). */ + readonly detailLoaded: boolean; +}): { + readonly messageIds: ReadonlyArray; + readonly source: "server" | "local" | "empty"; + readonly snapshotMissing: boolean; +} { + if (input.serverQueued.length > 0) { + // Dedupe while preserving server order. + const seen = new Set(); + const messageIds: MessageId[] = []; + for (const entry of input.serverQueued) { + const key = String(entry.messageId); + if (seen.has(key)) continue; + seen.add(key); + messageIds.push(entry.messageId); + } + return { messageIds, source: "server", snapshotMissing: false }; + } + + if (input.localPending.length > 0) { + const seen = new Set(); + const messageIds: MessageId[] = []; + for (const entry of input.localPending) { + const key = String(entry.t3MessageId); + if (seen.has(key)) continue; + seen.add(key); + messageIds.push(entry.t3MessageId); + } + return { + messageIds, + source: "local", + snapshotMissing: !input.detailLoaded, + }; + } + + return { + messageIds: [], + source: "empty", + snapshotMissing: !input.detailLoaded, + }; +} + +/** User-facing reply when steernow has nothing to inject. */ +export function formatSteernowEmptyQueueMessage(input: { + readonly snapshotMissing: boolean; +}): string { + if (input.snapshotMissing) { + return [ + "Could not load the server queue (thread snapshot unavailable), and nothing is parked in this bot process.", + "Use `/agent steer prompt:…` (or `@Omegent --steer …`) to inject mid-turn, then try `/agent steernow` again if you park follow-ups.", + ].join(" "); + } + return [ + "Nothing is queued on this thread.", + "Mid-turn follow-ups park with 📥 by default — then `/agent steernow` flushes them.", + "To inject immediately, use `/agent steer prompt:…` or `@Omegent --steer …`.", + ].join(" "); +} + export function createDiscordQueuedPromptRegistry() { const byDiscordMessageId = new Map(); const byT3ThreadId = new Map>(); diff --git a/apps/discord-bot/src/features/MentionRouter.ts b/apps/discord-bot/src/features/MentionRouter.ts index 7ff754c4176..915168e0419 100644 --- a/apps/discord-bot/src/features/MentionRouter.ts +++ b/apps/discord-bot/src/features/MentionRouter.ts @@ -81,7 +81,9 @@ import { ensureChannelInfoPin } from "./ChannelInfoPin.ts"; import { makeDiscordThreadTurnCoordinator } from "./DiscordThreadTurnCoordinator.ts"; import { createDiscordQueuedPromptRegistry, + formatSteernowEmptyQueueMessage, QUEUED_PROMPT_REACTION_EMOJI, + resolveSteernowMessageIds, } from "./DiscordQueuedPromptRegistry.ts"; import { BridgeHub } from "./BridgeHub.ts"; import { bridgeThreadToDiscord, getLiveDiscordBridge } from "./ResponseBridge.ts"; @@ -941,7 +943,7 @@ const make = (botConfig: DiscordBotConfig) => const followUpDelivery = resolveDiscordFollowUpDelivery(input.flags); const t3MessageId = MessageId.make(startedTurn.messageId); if (followUpDelivery === "steer" && turnAlreadyRunning) { - yield* t3 + const steered = yield* t3 .steerQueuedMessage({ threadId: existing.t3ThreadId, messageId: t3MessageId, @@ -953,6 +955,7 @@ const make = (botConfig: DiscordBotConfig) => messageId: startedTurn.messageId, }), ), + Effect.as(true), Effect.catch((error) => Effect.logWarning( "Steer after startTurn failed; message may remain server-queued", @@ -961,9 +964,29 @@ const make = (botConfig: DiscordBotConfig) => messageId: startedTurn.messageId, error: String(error), }, - ), + ).pipe(Effect.as(false)), ), ); + // If steer raced pending turn-start (or similar), keep the Discord-side + // registry + 📥 badge so /omegent steernow can flush without waiting on + // a lagging HTTP thread snapshot. + if (!steered) { + const discordMessageId = + typeof input.mentionMessage?.id === "string" ? input.mentionMessage.id : null; + if (discordMessageId !== null) { + const authorUserId = + typeof input.mentionMessage?.author?.id === "string" + ? input.mentionMessage.author.id + : null; + yield* markDiscordPromptQueued({ + discordChannelId: input.discordThreadId, + discordMessageId, + t3ThreadId: existing.t3ThreadId, + t3MessageId, + authorUserId, + }); + } + } } else if (followUpDelivery === "queue" && turnAlreadyRunning) { const discordMessageId = typeof input.mentionMessage?.id === "string" ? input.mentionMessage.id : null; @@ -2835,32 +2858,51 @@ const make = (botConfig: DiscordBotConfig) => }); } + // Server snapshot is authoritative after restart; local registry covers + // HTTP lag / soft-failed fetchThreadDetail so just-queued 📥 items still flush. const detail = yield* t3.fetchThreadDetail(link.t3ThreadId); - const queued = detail?.thread.queuedMessages ?? []; - if (queued.length === 0) { - return slashReply("Nothing is queued on this thread.", { ephemeral: true }); + const resolved = resolveSteernowMessageIds({ + serverQueued: detail?.thread.queuedMessages ?? [], + localPending: queuedPrompts.listForThread(link.t3ThreadId), + detailLoaded: detail !== null, + }); + if (resolved.messageIds.length === 0) { + return slashReply( + formatSteernowEmptyQueueMessage({ + snapshotMissing: resolved.snapshotMissing, + }), + { ephemeral: true }, + ); + } + + if (resolved.source === "local") { + yield* Effect.logInfo("steernow using local queued-prompt registry fallback", { + t3ThreadId: link.t3ThreadId, + count: resolved.messageIds.length, + snapshotMissing: resolved.snapshotMissing, + }); } let steered = 0; const failures: string[] = []; - for (const message of queued) { + for (const messageId of resolved.messageIds) { const result = yield* t3 .steerQueuedMessage({ threadId: link.t3ThreadId, - messageId: message.messageId, + messageId, }) .pipe(Effect.result); if (Result.isSuccess(result)) { steered += 1; - const pending = queuedPrompts.forgetT3Message(link.t3ThreadId, message.messageId); + const pending = queuedPrompts.forgetT3Message(link.t3ThreadId, messageId); if (pending !== null) { yield* clearQueuedBadge(pending); } } else { - failures.push(String(message.messageId)); + failures.push(String(messageId)); yield* Effect.logWarning("steernow failed for queued message", { t3ThreadId: link.t3ThreadId, - messageId: message.messageId, + messageId, error: String(result.failure), }); } @@ -2868,7 +2910,7 @@ const make = (botConfig: DiscordBotConfig) => if (steered === 0) { return slashReply( - `Could not steer the queue (${failures.length} failed). Try again when the turn is running.`, + `Could not steer the queue (${failures.length} failed). Try again when the turn is running, or use \`/agent steer prompt:…\` to inject a new mid-turn prompt.`, { ephemeral: true }, ); }