Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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", () => {
Expand Down Expand Up @@ -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");
});
});
70 changes: 70 additions & 0 deletions apps/discord-bot/src/features/DiscordQueuedPromptRegistry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<MessageId>;
readonly source: "server" | "local" | "empty";
readonly snapshotMissing: boolean;
} {
if (input.serverQueued.length > 0) {
// Dedupe while preserving server order.
const seen = new Set<string>();
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<string>();
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<string, PendingQueuedPrompt>();
const byT3ThreadId = new Map<string, Set<string>>();
Expand Down
64 changes: 53 additions & 11 deletions apps/discord-bot/src/features/MentionRouter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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,
Expand All @@ -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",
Expand All @@ -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;
Expand Down Expand Up @@ -2835,40 +2858,59 @@ 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),
});
}
}

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 },
);
}
Expand Down