diff --git a/PRODUCT.md b/PRODUCT.md index ab04d1834..74262fa0c 100644 --- a/PRODUCT.md +++ b/PRODUCT.md @@ -53,7 +53,7 @@ No public inbound ports are required for normal usage. - Send text prompts to OpenCode - Accept voice/audio messages, transcribe via Whisper-compatible STT API, and forward recognized text as prompts - Interrupt current task (ESC equivalent) -- Optionally queue text, transcribed voice, photos, rich formatted messages with photos, supported documents, and media groups sent while a task is running; hold at most `MAX_QUEUED_PROMPTS` (5) items and 20 MiB of raw Telegram media bytes, checked from reliable `file_size` before downloads +- Optionally admit text, transcribed voice, photos, rich formatted messages with photos, supported documents, and media groups to OpenCode's ordered session queue while a task is running; hold at most `MAX_QUEUED_PROMPTS` (5) pending bot admissions and 20 MiB of raw Telegram media bytes, checked from reliable `file_size` before downloads - Handle OpenCode questions with inline options and custom text answers - Send selected/custom answers back to OpenCode (`question.reply`) - Handle permission requests interactively (`allow once` / `always` / `reject`) diff --git a/src/app/managers/external-input-suppression-manager.ts b/src/app/managers/external-input-suppression-manager.ts index 288e6ef1e..98e6d3508 100644 --- a/src/app/managers/external-input-suppression-manager.ts +++ b/src/app/managers/external-input-suppression-manager.ts @@ -1,74 +1,106 @@ -const SUPPRESSION_TTL_MS = 60_000; - -interface SuppressionEntry { - text: string; - createdAt: number; -} - -function normalizeExternalUserInputText(text: string): string { - return text.replace(/\r\n/g, "\n").trim(); -} +const MESSAGE_SUPPRESSION_TTL_MS = 5 * 60_000; class ExternalUserInputSuppressionManager { - private entriesBySession = new Map(); + private messageIdsBySession = new Map>(); + private messagePruneTimer: ReturnType | null = null; - register(sessionId: string, text: string, now: number = Date.now()): void { - const normalizedText = normalizeExternalUserInputText(text); - if (!sessionId || !normalizedText) { + registerMessage(sessionId: string, messageId: string, now: number = Date.now()): void { + if (!sessionId || !messageId) { return; } - this.prune(now); - - const sessionEntries = this.entriesBySession.get(sessionId) ?? []; - sessionEntries.push({ text: normalizedText, createdAt: now }); - this.entriesBySession.set(sessionId, sessionEntries); + this.pruneMessages(now); + const messageIds = this.messageIdsBySession.get(sessionId) ?? new Map(); + messageIds.set(messageId, now + MESSAGE_SUPPRESSION_TTL_MS); + this.messageIdsBySession.set(sessionId, messageIds); + this.scheduleMessagePrune(now); } - consume(sessionId: string, text: string, now: number = Date.now()): boolean { - const normalizedText = normalizeExternalUserInputText(text); - if (!sessionId || !normalizedText) { - return false; + consumeMessage(sessionId: string, messageId: string): boolean { + this.pruneMessages(Date.now()); + const messageIds = this.messageIdsBySession.get(sessionId); + if (messageIds?.delete(messageId)) { + if (messageIds.size === 0) { + this.messageIdsBySession.delete(sessionId); + } + this.scheduleMessagePrune(Date.now()); + return true; } - this.prune(now); - - const sessionEntries = this.entriesBySession.get(sessionId); - if (!sessionEntries?.length) { - return false; - } + return false; + } - const entryIndex = sessionEntries.findIndex((entry) => entry.text === normalizedText); - if (entryIndex < 0) { - return false; + discardMessage(sessionId: string, messageId: string): void { + const messageIds = this.messageIdsBySession.get(sessionId); + if (!messageIds) { + return; } - sessionEntries.splice(entryIndex, 1); - if (sessionEntries.length === 0) { - this.entriesBySession.delete(sessionId); + messageIds.delete(messageId); + if (messageIds.size === 0) { + this.messageIdsBySession.delete(sessionId); } + this.scheduleMessagePrune(Date.now()); + } - return true; + clearSession(sessionId: string): void { + this.messageIdsBySession.delete(sessionId); + this.scheduleMessagePrune(Date.now()); } clearAll(): void { - this.entriesBySession.clear(); + this.messageIdsBySession.clear(); + if (this.messagePruneTimer) { + clearTimeout(this.messagePruneTimer); + this.messagePruneTimer = null; + } } __resetForTests(): void { this.clearAll(); } - private prune(now: number): void { - for (const [sessionId, sessionEntries] of this.entriesBySession.entries()) { - const activeEntries = sessionEntries.filter((entry) => now - entry.createdAt <= SUPPRESSION_TTL_MS); - if (activeEntries.length === 0) { - this.entriesBySession.delete(sessionId); - continue; + __getMessageCountForTests(): number { + return Array.from(this.messageIdsBySession.values()).reduce( + (count, messageIds) => count + messageIds.size, + 0, + ); + } + + private pruneMessages(now: number): void { + for (const [sessionId, messageIds] of this.messageIdsBySession.entries()) { + for (const [messageId, expiresAt] of messageIds.entries()) { + if (expiresAt <= now) { + messageIds.delete(messageId); + } } + if (messageIds.size === 0) { + this.messageIdsBySession.delete(sessionId); + } + } + } - this.entriesBySession.set(sessionId, activeEntries); + private scheduleMessagePrune(now: number): void { + if (this.messagePruneTimer) { + clearTimeout(this.messagePruneTimer); + this.messagePruneTimer = null; } + + const expirations = Array.from(this.messageIdsBySession.values()).flatMap((messageIds) => + Array.from(messageIds.values()), + ); + if (expirations.length === 0) { + return; + } + + const expiresAt = Math.min(...expirations); + this.messagePruneTimer = setTimeout(() => { + this.messagePruneTimer = null; + const pruneAt = Date.now(); + this.pruneMessages(pruneAt); + this.scheduleMessagePrune(pruneAt); + }, Math.max(0, expiresAt - now)); + this.messagePruneTimer.unref?.(); } } diff --git a/src/app/managers/interaction-manager.ts b/src/app/managers/interaction-manager.ts index 9fbb16648..798333ead 100644 --- a/src/app/managers/interaction-manager.ts +++ b/src/app/managers/interaction-manager.ts @@ -91,6 +91,13 @@ class InteractionManager { return cloneState(nextState); } + tryStart(options: StartInteractionOptions): InteractionState | null { + if (this.state) { + return null; + } + return this.start(options); + } + get(): InteractionState | null { if (!this.state) { return null; @@ -146,6 +153,25 @@ class InteractionManager { return cloneState(this.state); } + transitionIfMetadata( + key: string, + value: unknown, + options: TransitionInteractionOptions, + ): InteractionState | null { + if (!this.state || this.state.metadata[key] !== value) { + return null; + } + return this.transition(options); + } + + clearIfMetadata(key: string, value: unknown, reason: InteractionClearReason): boolean { + if (!this.state || this.state.metadata[key] !== value) { + return false; + } + this.clear(reason); + return true; + } + clear(reason: InteractionClearReason = "manual"): void { if (!this.state) { return; diff --git a/src/app/managers/prompt-response-mode-manager.ts b/src/app/managers/prompt-response-mode-manager.ts new file mode 100644 index 000000000..61933243b --- /dev/null +++ b/src/app/managers/prompt-response-mode-manager.ts @@ -0,0 +1,35 @@ +export type PromptResponseMode = "text_only" | "text_and_tts"; + +const promptResponseModes = new Map>(); + +export function setPromptResponseMode( + sessionId: string, + messageId: string, + responseMode: PromptResponseMode, +): void { + const modes = promptResponseModes.get(sessionId) ?? new Map(); + modes.set(messageId, responseMode); + promptResponseModes.set(sessionId, modes); +} + +export function clearPromptResponseMode(sessionId: string, messageId?: string): void { + if (!messageId) { + promptResponseModes.delete(sessionId); + return; + } + + const modes = promptResponseModes.get(sessionId); + modes?.delete(messageId); + if (modes?.size === 0) { + promptResponseModes.delete(sessionId); + } +} + +export function consumePromptResponseMode( + sessionId: string, + messageId: string, +): PromptResponseMode | null { + const responseMode = promptResponseModes.get(sessionId)?.get(messageId) ?? null; + clearPromptResponseMode(sessionId, messageId); + return responseMode; +} diff --git a/src/app/managers/summary-aggregation-manager.ts b/src/app/managers/summary-aggregation-manager.ts index 1db68398b..790b2930a 100644 --- a/src/app/managers/summary-aggregation-manager.ts +++ b/src/app/managers/summary-aggregation-manager.ts @@ -18,6 +18,7 @@ export interface SummaryInfo { } export interface MessageCompletionInfo { + parentMessageId?: string | undefined; agent?: string | undefined; providerID?: string | undefined; modelID?: string | undefined; @@ -1266,6 +1267,7 @@ class SummaryAggregator { if (this.onCompleteCallback && finalText.length > 0) { this.onCompleteCallback(this.currentSessionId!, messageID, finalText, { + parentMessageId: info.parentID, agent: info.agent, providerID: info.providerID, modelID: info.modelID, diff --git a/src/app/managers/telegram-input-order-manager.ts b/src/app/managers/telegram-input-order-manager.ts new file mode 100644 index 000000000..041f9e04d --- /dev/null +++ b/src/app/managers/telegram-input-order-manager.ts @@ -0,0 +1,103 @@ +import type { Context, NextFunction } from "grammy"; + +interface PendingWaiter { + messageId: number; + resolve: () => void; +} + +class TelegramInputOrderManager { + private readonly pendingByChat = new Map>(); + private readonly waitersByChat = new Map(); + + defer(chatId: number, messageId: number): void { + const pending = this.pendingByChat.get(chatId) ?? new Set(); + pending.add(messageId); + this.pendingByChat.set(chatId, pending); + } + + release(chatId: number, messageId: number): void { + const pending = this.pendingByChat.get(chatId); + pending?.delete(messageId); + if (pending?.size === 0) { + this.pendingByChat.delete(chatId); + } + this.resolveReadyWaiters(chatId); + } + + waitForEarlier(chatId: number, messageId: number): Promise { + if (!this.hasEarlierPending(chatId, messageId)) { + return Promise.resolve(); + } + + return new Promise((resolve) => { + const waiters = this.waitersByChat.get(chatId) ?? []; + waiters.push({ messageId, resolve }); + this.waitersByChat.set(chatId, waiters); + }); + } + + __resetForTests(): void { + for (const waiters of this.waitersByChat.values()) { + for (const waiter of waiters) { + waiter.resolve(); + } + } + this.pendingByChat.clear(); + this.waitersByChat.clear(); + } + + private hasEarlierPending(chatId: number, messageId: number): boolean { + return Array.from(this.pendingByChat.get(chatId) ?? []).some( + (pendingMessageId) => pendingMessageId < messageId, + ); + } + + private resolveReadyWaiters(chatId: number): void { + const waiters = this.waitersByChat.get(chatId); + if (!waiters) { + return; + } + + const blocked: PendingWaiter[] = []; + for (const waiter of waiters) { + if (this.hasEarlierPending(chatId, waiter.messageId)) { + blocked.push(waiter); + } else { + waiter.resolve(); + } + } + + if (blocked.length > 0) { + this.waitersByChat.set(chatId, blocked); + } else { + this.waitersByChat.delete(chatId); + } + } +} + +export const telegramInputOrderManager = new TelegramInputOrderManager(); + +export async function telegramInputOrderMiddleware( + ctx: Context, + next: NextFunction, +): Promise { + const message = ctx.message; + const chatId = ctx.chat?.id; + if (chatId === undefined) { + await next(); + return; + } + + if (!message) { + await next(); + return; + } + + if (message.media_group_id) { + await next(); + return; + } + + await telegramInputOrderManager.waitForEarlier(chatId, message.message_id); + await next(); +} diff --git a/src/app/services/external-user-input-service.ts b/src/app/services/external-user-input-service.ts index 53f360d17..2ab9cfe71 100644 --- a/src/app/services/external-user-input-service.ts +++ b/src/app/services/external-user-input-service.ts @@ -8,6 +8,11 @@ export interface ExternalUserInputNotification { rawFallbackText: string; } +export type ConsumeSuppressedInput = ( + sessionId: string, + messageId: string, +) => boolean; + function normalizeExternalUserInputText(text: string): string { return text.replace(/\r\n/g, "\n").trim(); } diff --git a/src/app/services/session-service.ts b/src/app/services/session-service.ts index da45a617d..b1a13e16a 100644 --- a/src/app/services/session-service.ts +++ b/src/app/services/session-service.ts @@ -6,15 +6,22 @@ import { import { promptQueue } from "../managers/prompt-queue-manager.js"; import { promptAttachment } from "../managers/prompt-attachment-manager.js"; import type { SessionInfo } from "../types/session.js"; +import { externalUserInputSuppressionManager } from "../managers/external-input-suppression-manager.js"; +import { clearPromptResponseMode } from "../managers/prompt-response-mode-manager.js"; export type { SessionInfo }; export function setCurrentSession(sessionInfo: SessionInfo): void { // Renaming reuses this setter with the same id, so only an actual session // switch may drop prompts queued for the previous session. - if (getSettingsSession()?.id !== sessionInfo.id) { + const previousSessionId = getSettingsSession()?.id; + if (previousSessionId !== sessionInfo.id) { promptQueue.clear("session_switched"); promptAttachment.clear("session_switched"); + if (previousSessionId) { + externalUserInputSuppressionManager.clearSession(previousSessionId); + clearPromptResponseMode(previousSessionId); + } } setSettingsSession(sessionInfo); @@ -25,7 +32,12 @@ export function getCurrentSession(): SessionInfo | null { } export function clearSession(): void { + const previousSessionId = getSettingsSession()?.id; promptQueue.clear("session_cleared"); promptAttachment.clear("session_cleared"); + if (previousSessionId) { + externalUserInputSuppressionManager.clearSession(previousSessionId); + clearPromptResponseMode(previousSessionId); + } clearSettingsSession(); } diff --git a/src/app/services/tts-service.ts b/src/app/services/tts-service.ts index 6804677a4..d6ad54451 100644 --- a/src/app/services/tts-service.ts +++ b/src/app/services/tts-service.ts @@ -16,8 +16,9 @@ type PromptResponseMode = "text_only" | "text_and_tts"; interface PrepareTtsResponseParams { sessionId: string; + promptMessageId: string; text: string; - consumeResponseMode: (sessionId: string) => PromptResponseMode | null; + consumeResponseMode: (sessionId: string, messageId: string) => PromptResponseMode | null; isTtsConfigured?: (() => boolean) | undefined; synthesizeSpeech?: ((text: string) => Promise) | undefined; } @@ -273,12 +274,13 @@ export async function synthesizeSpeech(text: string): Promise { export async function prepareTtsResponseForSession({ sessionId, + promptMessageId, text, consumeResponseMode, isTtsConfigured: isTtsConfiguredImpl = isTtsConfigured, synthesizeSpeech: synthesizeSpeechImpl = synthesizeSpeech, }: PrepareTtsResponseParams): Promise { - const responseMode = consumeResponseMode(sessionId); + const responseMode = consumeResponseMode(sessionId, promptMessageId); if (responseMode !== "text_and_tts") { return { shouldSend: false }; } diff --git a/src/bot/callbacks/command-catalog-callback-handler.ts b/src/bot/callbacks/command-catalog-callback-handler.ts index 2dbdf7b31..1683f8a55 100644 --- a/src/bot/callbacks/command-catalog-callback-handler.ts +++ b/src/bot/callbacks/command-catalog-callback-handler.ts @@ -27,6 +27,7 @@ import { markAttachedSessionIdle, } from "../../app/services/attach-service.js"; import { externalUserInputSuppressionManager } from "../../app/managers/external-input-suppression-manager.js"; +import { createOpencodeMessageId } from "../../utils/opencode-message-id.js"; import { opencodeClient } from "../../opencode/client.js"; import { buildCommandsConfirmKeyboard, @@ -285,16 +286,15 @@ export async function executeCommand( configuredProviderID: storedModel.providerID, configuredModelID: storedModel.modelID, }); - externalUserInputSuppressionManager.register( - session.id, - args ? `/${params.commandName} ${args}` : `/${params.commandName}`, - ); + const promptMessageId = createOpencodeMessageId(); + externalUserInputSuppressionManager.registerMessage(session.id, promptMessageId); safeBackgroundTask({ taskName: "session.command", task: () => opencodeClient.session.command({ sessionID: session.id, + messageID: promptMessageId, directory: session.directory, command: params.commandName, arguments: args, @@ -307,6 +307,7 @@ export async function executeCommand( foregroundSessionState.markIdle(session.id); void markAttachedSessionIdle(session.id); assistantRunState.clearRun(session.id, "session_command_api_error"); + externalUserInputSuppressionManager.discardMessage(session.id, promptMessageId); logger.error("[Commands] OpenCode API returned an error for session.command", { sessionId: session.id, command: params.commandName, diff --git a/src/bot/handlers/media-group-handler.ts b/src/bot/handlers/media-group-handler.ts index 6ab434710..8f4972dfb 100644 --- a/src/bot/handlers/media-group-handler.ts +++ b/src/bot/handlers/media-group-handler.ts @@ -19,6 +19,10 @@ import { rejectQueuedMediaBeforePreparation, tryEnqueuePromptIfBusy, } from "./prompt-queue-dispatch.js"; +import { telegramInputOrderManager } from "../../app/managers/telegram-input-order-manager.js"; +import { getCurrentSession } from "../../app/services/session-service.js"; +import { getCurrentProject } from "../../app/stores/settings-store.js"; +import type { PromptTarget } from "./prompt.js"; const DEFAULT_MEDIA_GROUP_DEBOUNCE_MS = 1_000; @@ -62,6 +66,7 @@ type ValidMediaGroupItem = interface MediaGroupBatch { timer: ReturnType; items: PendingMediaGroupItem[]; + target?: PromptTarget; } interface ValidatedMediaGroup { @@ -86,6 +91,7 @@ export interface MediaGroupHandlerDeps extends ProcessPromptDeps { ctx: Context, input: IncomingPrompt, deps: ProcessPromptDeps, + options?: { target?: PromptTarget }, ) => Promise; } @@ -121,6 +127,7 @@ export class MediaGroupAttachmentHandler { flushPendingPrompt(chatId); const key = this.getBatchKey(chatId, mediaGroupId); + telegramInputOrderManager.defer(chatId, item.messageId); const existingBatch = this.batches.get(key); if (existingBatch) { @@ -130,9 +137,11 @@ export class MediaGroupAttachmentHandler { return; } + const target = this.captureTarget(); this.batches.set(key, { items: [item], timer: this.createFlushTimer(key), + ...(target ? { target } : {}), }); } @@ -206,6 +215,12 @@ export class MediaGroupAttachmentHandler { logger.info(`[MediaGroup] Processing Telegram media group: key=${key}, items=${items.length}`); try { + const chatId = replyCtx.chat?.id; + const firstMessageId = items[0]?.messageId; + if (chatId !== undefined && firstMessageId !== undefined) { + await telegramInputOrderManager.waitForEarlier(chatId, firstMessageId); + } + const unsupportedContexts = items .filter((item) => item.kind === "unsupported") .map((item) => item.ctx); @@ -255,19 +270,42 @@ export class MediaGroupAttachmentHandler { await tryEnqueuePromptIfBusy(replyCtx, { ...createIncomingPrompt(promptText, { fileParts }), displayText: captions.join(" / ") || `[Album: ${items.length} files]`, - fileParts, ...(mediaBytes === undefined ? {} : { mediaBytes }), + ...(batch.target?.sessionId + ? { sessionId: batch.target.sessionId, directory: batch.target.directory } + : {}), }) ) { return; } - await processPrompt(replyCtx, createIncomingPrompt(promptText, { fileParts }), this.deps); + await processPrompt(replyCtx, createIncomingPrompt(promptText, { fileParts }), this.deps, { + ...(batch.target ? { target: batch.target } : {}), + }); } catch (err) { logger.error(`[MediaGroup] Failed to process media group: key=${key}`, err); await replyCtx.reply(t("bot.media_group_download_error")); + } finally { + const chatId = replyCtx.chat?.id; + if (chatId !== undefined) { + for (const item of items) { + telegramInputOrderManager.release(chatId, item.messageId); + } + } } } + private captureTarget(): PromptTarget | undefined { + const project = getCurrentProject(); + if (!project) { + return undefined; + } + + return { + sessionId: getCurrentSession()?.id ?? null, + directory: project.worktree, + }; + } + private async validateItems( items: PendingMediaGroupItem[], ): Promise { diff --git a/src/bot/handlers/prompt.ts b/src/bot/handlers/prompt.ts index 03d8242a9..e0195aeeb 100644 --- a/src/bot/handlers/prompt.ts +++ b/src/bot/handlers/prompt.ts @@ -1,6 +1,7 @@ import { Bot, Context } from "grammy"; import type { FilePartInput, TextPartInput } from "@opencode-ai/sdk/v2"; import type { Model } from "@opencode-ai/sdk/v2"; +import { createOpencodeMessageId } from "../../utils/opencode-message-id.js"; import { opencodeClient } from "../../opencode/client.js"; import { clearSession, @@ -44,16 +45,31 @@ import { supportsInput, } from "../../app/services/model-capabilities-service.js"; import type { IncomingPrompt } from "../../app/types/prompt.js"; +import { + consumePromptResponseMode, + setPromptResponseMode, + type PromptResponseMode, +} from "../../app/managers/prompt-response-mode-manager.js"; + +export { + clearPromptResponseMode, + consumePromptResponseMode, + setPromptResponseMode, + type PromptResponseMode, +} from "../../app/managers/prompt-response-mode-manager.js"; /** Module-level references for async callbacks that don't have ctx. */ let botInstance: Bot | null = null; let chatIdInstance: number | null = null; -const promptResponseModes = new Map(); -export type PromptResponseMode = "text_only" | "text_and_tts"; +export type PromptTarget = { + sessionId: string | null; + directory: string; +}; type ProcessPromptOptions = { responseMode?: PromptResponseMode; + target?: PromptTarget; }; export function getPromptBotInstance(): Bot | null { @@ -64,18 +80,9 @@ export function getPromptChatId(): number | null { return chatIdInstance; } -export function setPromptResponseMode(sessionId: string, responseMode: PromptResponseMode): void { - promptResponseModes.set(sessionId, responseMode); -} - -export function clearPromptResponseMode(sessionId: string): void { - promptResponseModes.delete(sessionId); -} - -export function consumePromptResponseMode(sessionId: string): PromptResponseMode | null { - const responseMode = promptResponseModes.get(sessionId) ?? null; - promptResponseModes.delete(sessionId); - return responseMode; +export function discardPromptDeliveryState(sessionId: string, messageId: string): void { + consumePromptResponseMode(sessionId, messageId); + externalUserInputSuppressionManager.discardMessage(sessionId, messageId); } async function isSessionBusy(sessionId: string, directory: string): Promise { @@ -193,6 +200,17 @@ export async function processUserPrompt( let currentSession = getCurrentSession(); let createdNewSession = false; + if ( + options.target && + (currentProject.worktree !== options.target.directory || + (currentSession?.id ?? null) !== options.target.sessionId || + (currentSession && currentSession.directory !== options.target.directory)) + ) { + logger.warn("[Bot] Refusing prompt after captured context changed"); + await ctx.reply(t("bot.prompt_send_error")); + return false; + } + if (currentSession && currentSession.directory !== currentProject.worktree) { logger.warn( `[Bot] Session/project mismatch detected. sessionDirectory=${currentSession.directory}, projectDirectory=${currentProject.worktree}. Resetting session context.`, @@ -369,11 +387,9 @@ export async function processUserPrompt( configuredProviderID: storedModel.providerID, configuredModelID: storedModel.modelID, }); - setPromptResponseMode(currentSession.id, responseMode); - - if (preparedInput.text.trim().length > 0) { - externalUserInputSuppressionManager.register(currentSession.id, preparedInput.text); - } + const promptMessageId = createOpencodeMessageId(); + setPromptResponseMode(currentSession.id, promptMessageId, responseMode); + externalUserInputSuppressionManager.registerMessage(currentSession.id, promptMessageId); // CRITICAL: Use the async prompt start endpoint here. // session.prompt streams the full assistant response and can outlive the original @@ -382,13 +398,14 @@ export async function processUserPrompt( // The actual assistant result still arrives via the SSE event subscription. safeBackgroundTask({ taskName: "session.promptAsync", - task: () => opencodeClient.session.promptAsync(promptOptions), + task: () => + opencodeClient.session.promptAsync({ ...promptOptions, messageID: promptMessageId }), onSuccess: ({ error }) => { if (error) { foregroundSessionState.markIdle(currentSession.id); void markAttachedSessionIdle(currentSession.id); assistantRunState.clearRun(currentSession.id, "session_prompt_api_error"); - clearPromptResponseMode(currentSession.id); + discardPromptDeliveryState(currentSession.id, promptMessageId); const details = formatErrorDetails(error, 6000); logger.error( "[Bot] OpenCode API returned an error for session.promptAsync", @@ -407,10 +424,6 @@ export async function processUserPrompt( logger.info("[Bot] session.promptAsync accepted"); }, onError: (error) => { - foregroundSessionState.markIdle(currentSession.id); - void markAttachedSessionIdle(currentSession.id); - assistantRunState.clearRun(currentSession.id, "session_prompt_background_error"); - clearPromptResponseMode(currentSession.id); const details = formatErrorDetails(error, 6000); logger.error("[Bot] session.promptAsync background task failed", promptErrorLogContext); logger.error("[Bot] session.promptAsync background failure details:", details); diff --git a/src/bot/handlers/tts-response-handler.ts b/src/bot/handlers/tts-response-handler.ts index 61b14afd9..cc29ec9d1 100644 --- a/src/bot/handlers/tts-response-handler.ts +++ b/src/bot/handlers/tts-response-handler.ts @@ -12,9 +12,13 @@ interface TelegramAudioApi { interface SendTtsResponseParams { api: TelegramAudioApi; sessionId: string; + promptMessageId?: string | undefined; chatId: number; text: string; - consumeResponseMode?: (sessionId: string) => "text_only" | "text_and_tts" | null; + consumeResponseMode?: ( + sessionId: string, + messageId: string, + ) => "text_only" | "text_and_tts" | null; isTtsConfigured?: () => boolean; synthesizeSpeech?: (text: string) => Promise; } @@ -22,15 +26,21 @@ interface SendTtsResponseParams { export async function sendTtsResponseForSession({ api, sessionId, + promptMessageId, chatId, text, consumeResponseMode: consumeResponseModeImpl = consumePromptResponseMode, isTtsConfigured, synthesizeSpeech, }: SendTtsResponseParams): Promise { + if (!promptMessageId) { + return false; + } + try { const prepared = await prepareTtsResponseForSession({ sessionId, + promptMessageId, text, consumeResponseMode: consumeResponseModeImpl, isTtsConfigured, diff --git a/src/bot/index.ts b/src/bot/index.ts index e60131f18..b3803a869 100644 --- a/src/bot/index.ts +++ b/src/bot/index.ts @@ -17,6 +17,7 @@ import { initializePromptQueueDispatch } from "./handlers/prompt-queue-dispatch. import { normalizeRichMessage } from "./handlers/rich-message-handler.js"; import { authMiddleware } from "./middleware/auth.js"; import { interactionGuardMiddleware } from "./middleware/interaction-guard.js"; +import { telegramInputOrderMiddleware } from "../app/managers/telegram-input-order-manager.js"; import { staleUpdateMiddleware } from "./middleware/stale-update.js"; import { ensureCommandsInitialized, @@ -177,6 +178,7 @@ export function createBot(localCommandRegistry = LocalCommandRegistry.empty()): bot.use(staleUpdateMiddleware); bot.on("message:rich_message", normalizeRichMessage); bot.use((ctx, next) => ensureCommandsInitialized(ctx, next, localCommandRegistry)); + bot.use(telegramInputOrderMiddleware); bot.use((ctx, next) => interactionGuardMiddleware(ctx, next, localCommandRegistry)); registerCommandRouter(bot, { diff --git a/src/bot/menus/inline-menu.ts b/src/bot/menus/inline-menu.ts index f09144f88..028747cf8 100644 --- a/src/bot/menus/inline-menu.ts +++ b/src/bot/menus/inline-menu.ts @@ -1,4 +1,5 @@ import { Context, InlineKeyboard } from "grammy"; +import { randomUUID } from "node:crypto"; import { interactionManager } from "../../app/managers/interaction-manager.js"; import type { InteractionMetadata, InteractionState } from "../../app/types/interaction.js"; import { logger } from "../../utils/logger.js"; @@ -101,8 +102,27 @@ export function appendInlineMenuCancelButton( export async function replyWithInlineMenu( ctx: Context, options: InlineMenuReplyOptions, -): Promise { +): Promise { const keyboard = appendInlineMenuCancelButton(options.keyboard, options.menuKind); + const activeMenu = getActiveInlineMenuMetadata(interactionManager.getSnapshot()); + if (options.menuKind === "settings" && activeMenu?.menuKind === "settings") { + interactionManager.clear("settings_menu_reopened"); + } + const reservationId = randomUUID(); + const reserved = interactionManager.tryStart({ + kind: "inline", + expectedInput: "callback", + metadata: { + ...options.metadata, + menuKind: options.menuKind, + reservationId, + }, + }); + if (!reserved) { + await ctx.reply(t("interaction.blocked.finish_current")); + return null; + } + const replyOptions: { reply_markup: InlineKeyboard; parse_mode?: "Markdown" | "HTML"; @@ -114,17 +134,27 @@ export async function replyWithInlineMenu( replyOptions.parse_mode = options.parseMode; } - const message = await ctx.reply(options.text, replyOptions); + let message; + try { + message = await ctx.reply(options.text, replyOptions); + } catch (error) { + interactionManager.clearIfMetadata("reservationId", reservationId, "inline_send_failed"); + throw error; + } - interactionManager.start({ - kind: "inline", - expectedInput: "callback", + const finalized = interactionManager.transitionIfMetadata("reservationId", reservationId, { metadata: { ...options.metadata, menuKind: options.menuKind, messageId: message.message_id, }, }); + if (!finalized) { + if (ctx.chat) { + await ctx.api.editMessageReplyMarkup(ctx.chat.id, message.message_id).catch(() => {}); + } + return null; + } logger.debug( `[InlineMenu] Opened menu: kind=${options.menuKind}, messageId=${message.message_id}`, diff --git a/src/bot/messages/external-user-input-notification.ts b/src/bot/messages/external-user-input-notification.ts index 3dfe22f5f..78bd0b041 100644 --- a/src/bot/messages/external-user-input-notification.ts +++ b/src/bot/messages/external-user-input-notification.ts @@ -1,6 +1,7 @@ import type { Api, RawApi } from "grammy"; import { buildExternalUserInputNotification, + type ConsumeSuppressedInput, type ExternalUserInputNotification, } from "../../app/services/external-user-input-service.js"; import { sendBotText } from "./telegram-text.js"; @@ -12,8 +13,9 @@ interface DeliverExternalUserInputParams { chatId: number; currentSessionId: string | null; sessionId: string; + messageId: string; text: string; - consumeSuppressedInput: (sessionId: string, text: string) => boolean; + consumeSuppressedInput: ConsumeSuppressedInput; } async function sendExternalUserInputNotification( @@ -35,6 +37,7 @@ export async function deliverExternalUserInputNotification({ chatId, currentSessionId, sessionId, + messageId, text, consumeSuppressedInput, }: DeliverExternalUserInputParams): Promise { @@ -43,7 +46,7 @@ export async function deliverExternalUserInputNotification({ return false; } - if (consumeSuppressedInput(sessionId, text)) { + if (consumeSuppressedInput(sessionId, messageId)) { return false; } diff --git a/src/bot/middleware/interaction-guard-decision.ts b/src/bot/middleware/interaction-guard-decision.ts index ef7695493..2ce84425c 100644 --- a/src/bot/middleware/interaction-guard-decision.ts +++ b/src/bot/middleware/interaction-guard-decision.ts @@ -6,22 +6,52 @@ import type { GuardDecision, IncomingInputType, InteractionState, - InteractionKind, } from "../../app/types/interaction.js"; import { foregroundSessionState } from "../../app/managers/foreground-session-state-manager.js"; import { attachManager } from "../../app/managers/attach-manager.js"; import { QUEUED_PROMPT_BUTTON_TEXT_PATTERN } from "../message-patterns.js"; import type { LocalCommandRegistry } from "../../app/services/local-command-registry.js"; -const BUSY_ALLOWED_COMMANDS = ["/abort", "/detach", "/status", "/help", "/opencode_stop"] as const; +const BUSY_ALLOWED_COMMANDS = [ + "/abort", + "/detach", + "/status", + "/help", + "/opencode_stop", + "/settings", +] as const; const BUSY_ALLOWED_COMMAND_SET = new Set(BUSY_ALLOWED_COMMANDS); -function isBusyAllowedCommand(command: string | undefined, localCommandRegistry?: LocalCommandRegistry): boolean { - return Boolean(command && (BUSY_ALLOWED_COMMAND_SET.has(command) || localCommandRegistry?.allowsWhenBusy(command))); +function isBusyAllowedCommand( + command: string | undefined, + state: InteractionState | null, + localCommandRegistry?: LocalCommandRegistry, +): boolean { + if (!command) { + return false; + } + + if (localCommandRegistry?.allowsWhenBusy(command)) { + return true; + } + + if (!BUSY_ALLOWED_COMMAND_SET.has(command)) { + return false; + } + + return ( + command !== "/settings" || + !state || + (state.kind === "inline" && state.metadata.menuKind === "settings") + ); } -function allowsBusyInteraction(kind: InteractionKind | undefined): boolean { - return kind === "question" || kind === "permission"; +function allowsBusyInteraction(state: InteractionState): boolean { + return ( + state.kind === "question" || + state.kind === "permission" || + (state.kind === "inline" && state.metadata.menuKind === "settings") + ); } // Removing a queued prompt only makes sense while the session is busy, so the @@ -164,14 +194,14 @@ export function resolveInteractionGuardDecision( if (state && localCommandRegistry?.has(command)) { return createBusyBlockDecision(inputType, state, "command_not_allowed", command); } - if (isBusyAllowedCommand(command, localCommandRegistry)) { + if (isBusyAllowedCommand(command, state, localCommandRegistry)) { return createAllowDecision(inputType, state, command, true); } return createBusyBlockDecision(inputType, state, "command_not_allowed", command); } - if (state && allowsBusyInteraction(state.kind)) { + if (state && allowsBusyInteraction(state)) { if (state.expectedInput === "mixed") { if (inputType === "callback" || inputType === "text") { return createAllowDecision(inputType, state, command, true); diff --git a/src/bot/middleware/interaction-guard.ts b/src/bot/middleware/interaction-guard.ts index d2286fb6d..1ae96ff46 100644 --- a/src/bot/middleware/interaction-guard.ts +++ b/src/bot/middleware/interaction-guard.ts @@ -12,6 +12,7 @@ import { logger } from "../../utils/logger.js"; import { t } from "../../i18n/index.js"; import { getIncomingPrompt } from "../handlers/rich-message-handler.js"; import type { LocalCommandRegistry } from "../../app/services/local-command-registry.js"; +import { telegramInputOrderManager } from "../../app/managers/telegram-input-order-manager.js"; function getInteractionBlockedMessage( reason: BlockReason | undefined, @@ -143,6 +144,9 @@ export async function interactionGuardMiddleware( ); if (isQueueableInput && incomingPrompt) { + if (ctx.chat && ctx.message?.message_id !== undefined) { + await telegramInputOrderManager.waitForEarlier(ctx.chat.id, ctx.message.message_id); + } const mediaBytes = getQueuedPhotoMediaBytes(incomingPrompt); if (incomingPrompt.photos.length > 0) { if (await rejectQueuedMediaBeforePreparation(ctx, mediaBytes)) { diff --git a/src/bot/services/event-subscription-service.ts b/src/bot/services/event-subscription-service.ts index e2fe62caa..9a9f4f917 100644 --- a/src/bot/services/event-subscription-service.ts +++ b/src/bot/services/event-subscription-service.ts @@ -525,6 +525,7 @@ class EventSubscriptionService implements BotEventSubscriptionService { this.sessionCompletionTasks.clear(); this.clearToolElapsedState(null, reason); assistantRunState.clearAll(reason); + externalUserInputSuppressionManager.clearAll(); }; cleanup(reason: string): void { @@ -607,7 +608,9 @@ class EventSubscriptionService implements BotEventSubscriptionService { void this.enqueueSessionCompletionTask(sessionId, async () => { if (!this.botInstance || !this.chatIdInstance) { logger.error("Bot or chat ID not available for sending message"); - clearPromptResponseMode(sessionId); + if (completionInfo.parentMessageId) { + clearPromptResponseMode(sessionId, completionInfo.parentMessageId); + } this.clearAssistantResponseStream(sessionId, messageId, "bot_context_missing"); this.clearThinkingStream(sessionId, messageId, "bot_context_missing"); this.toolCallStreamer.clearSession(sessionId, "bot_context_missing"); @@ -620,7 +623,9 @@ class EventSubscriptionService implements BotEventSubscriptionService { const currentSession = getCurrentSession(); if (currentSession?.id !== sessionId) { - clearPromptResponseMode(sessionId); + if (completionInfo.parentMessageId) { + clearPromptResponseMode(sessionId, completionInfo.parentMessageId); + } this.clearAssistantResponseStream(sessionId, messageId, "session_mismatch"); this.clearThinkingStream(sessionId, messageId, "session_mismatch"); this.toolCallStreamer.clearSession(sessionId, "session_mismatch"); @@ -689,11 +694,14 @@ class EventSubscriptionService implements BotEventSubscriptionService { await sendTtsResponseForSession({ api: botApi, sessionId, + promptMessageId: completionInfo.parentMessageId, chatId, text: messageText, }); } catch (err) { - clearPromptResponseMode(sessionId); + if (completionInfo.parentMessageId) { + clearPromptResponseMode(sessionId, completionInfo.parentMessageId); + } this.clearThinkingStream(sessionId, messageId, "assistant_finalize_failed"); this.compactProgressStreamer.clearSession(sessionId, "assistant_finalize_failed"); assistantRunState.clearRun(sessionId, "assistant_finalize_failed"); @@ -707,7 +715,7 @@ class EventSubscriptionService implements BotEventSubscriptionService { }); }); - summaryAggregator.setOnExternalUserInput(async (sessionId, _messageId, messageText) => { + summaryAggregator.setOnExternalUserInput(async (sessionId, messageId, messageText) => { void this.enqueueSessionCompletionTask(sessionId, async () => { if (!this.botInstance || !this.chatIdInstance) { return; @@ -719,9 +727,13 @@ class EventSubscriptionService implements BotEventSubscriptionService { chatId: this.chatIdInstance, currentSessionId: getCurrentSession()?.id ?? null, sessionId, + messageId, text: messageText, - consumeSuppressedInput: (incomingSessionId, incomingText) => - externalUserInputSuppressionManager.consume(incomingSessionId, incomingText), + consumeSuppressedInput: (incomingSessionId, incomingMessageId) => + externalUserInputSuppressionManager.consumeMessage( + incomingSessionId, + incomingMessageId, + ), }); } catch (err) { logger.error("[Bot] Failed to deliver external user input to Telegram:", err); @@ -1163,7 +1175,6 @@ class EventSubscriptionService implements BotEventSubscriptionService { await this.sessionCompletionTasks.get(sessionId)?.catch(() => undefined); const completedRun = assistantRunState.finishRun(sessionId, "session_idle"); - clearPromptResponseMode(sessionId); if (!this.botInstance || !this.chatIdInstance) { foregroundSessionState.markIdle(sessionId); @@ -1218,7 +1229,6 @@ class EventSubscriptionService implements BotEventSubscriptionService { this.clearToolElapsedState(sessionId, "session_error"); if (!this.botInstance || !this.chatIdInstance) { - clearPromptResponseMode(sessionId); this.compactProgressStreamer.clearSession(sessionId, "session_error_no_bot_context"); assistantRunState.clearRun(sessionId, "session_error_no_bot_context"); foregroundSessionState.markIdle(sessionId); @@ -1227,8 +1237,7 @@ class EventSubscriptionService implements BotEventSubscriptionService { const currentSession = getCurrentSession(); if (!currentSession || currentSession.id !== sessionId) { - clearPromptResponseMode(sessionId); - this.clearAssistantResponseSession(sessionId, "session_error_not_current"); + this.clearAssistantResponseSession(sessionId, "session_error_not_current"); this.toolCallStreamer.clearSession(sessionId, "session_error_not_current"); this.compactProgressStreamer.clearSession(sessionId, "session_error_not_current"); assistantRunState.clearRun(sessionId, "session_error_not_current"); @@ -1239,7 +1248,6 @@ class EventSubscriptionService implements BotEventSubscriptionService { this.clearAssistantResponseSession(sessionId, "session_error"); this.compactProgressStreamer.clearSession(sessionId, "session_error"); - clearPromptResponseMode(sessionId); assistantRunState.clearRun(sessionId, "session_error"); await Promise.all([ this.toolMessageBatcher.flushSession(sessionId, "session_error"), diff --git a/src/utils/opencode-message-id.ts b/src/utils/opencode-message-id.ts new file mode 100644 index 000000000..6a3958c71 --- /dev/null +++ b/src/utils/opencode-message-id.ts @@ -0,0 +1,5 @@ +import { randomUUID } from "node:crypto"; + +export function createOpencodeMessageId(): string { + return `msg_${randomUUID()}`; +} diff --git a/tests/app/managers/external-input-suppression-manager.test.ts b/tests/app/managers/external-input-suppression-manager.test.ts index cae58110d..e7fcb41b7 100644 --- a/tests/app/managers/external-input-suppression-manager.test.ts +++ b/tests/app/managers/external-input-suppression-manager.test.ts @@ -1,4 +1,4 @@ -import { beforeEach, describe, expect, it } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { externalUserInputSuppressionManager } from "../../../src/app/managers/external-input-suppression-manager.js"; describe("external-input/suppression", () => { @@ -6,30 +6,62 @@ describe("external-input/suppression", () => { externalUserInputSuppressionManager.__resetForTests(); }); - it("consumes a matching suppressed input for the same session", () => { - externalUserInputSuppressionManager.register("session-1", "Review README"); + afterEach(() => { + vi.useRealTimers(); + }); + + it("retires an unobserved admitted identity after the bounded lifetime", () => { + vi.useFakeTimers(); + externalUserInputSuppressionManager.registerMessage("session-1", "message-1"); + expect(externalUserInputSuppressionManager.__getMessageCountForTests()).toBe(1); + + vi.advanceTimersByTime(5 * 60 * 1_000); + + expect(externalUserInputSuppressionManager.__getMessageCountForTests()).toBe(0); - expect(externalUserInputSuppressionManager.consume("session-1", "Review README")).toBe(true); - expect(externalUserInputSuppressionManager.consume("session-1", "Review README")).toBe(false); + expect( + externalUserInputSuppressionManager.consumeMessage("session-1", "message-1"), + ).toBe(false); }); - it("does not consume a suppressed input from another session", () => { - externalUserInputSuppressionManager.register("session-1", "Review README"); + it("refreshes the same stable identity when an ambiguous admission is retried", () => { + vi.useFakeTimers(); + externalUserInputSuppressionManager.registerMessage("session-1", "message-1"); + vi.advanceTimersByTime(4 * 60 * 1_000); + externalUserInputSuppressionManager.registerMessage("session-1", "message-1"); + vi.advanceTimersByTime(2 * 60 * 1_000); - expect(externalUserInputSuppressionManager.consume("session-2", "Review README")).toBe(false); + expect( + externalUserInputSuppressionManager.consumeMessage("session-1", "message-1"), + ).toBe(true); + expect(externalUserInputSuppressionManager.__getMessageCountForTests()).toBe(0); }); - it("does not consume different text", () => { - externalUserInputSuppressionManager.register("session-1", "Review README"); + it("does not suppress another message with the same text", () => { + externalUserInputSuppressionManager.registerMessage("session-1", "message-1"); - expect(externalUserInputSuppressionManager.consume("session-1", "Review tests")).toBe(false); + expect( + externalUserInputSuppressionManager.consumeMessage("session-1", "external-message"), + ).toBe(false); }); - it("expires stale suppression entries", () => { - externalUserInputSuppressionManager.register("session-1", "Review README", 1_000); + it("does not suppress the same identity in another session", () => { + externalUserInputSuppressionManager.registerMessage("session-1", "message-1"); + + expect( + externalUserInputSuppressionManager.consumeMessage("session-2", "message-1"), + ).toBe(false); + expect( + externalUserInputSuppressionManager.consumeMessage("session-1", "message-1"), + ).toBe(true); + }); + + it("returns identity storage to baseline when a session retires", () => { + externalUserInputSuppressionManager.registerMessage("session-1", "message-1"); + externalUserInputSuppressionManager.registerMessage("session-1", "message-2"); + + externalUserInputSuppressionManager.clearSession("session-1"); - expect(externalUserInputSuppressionManager.consume("session-1", "Review README", 61_001)).toBe( - false, - ); + expect(externalUserInputSuppressionManager.__getMessageCountForTests()).toBe(0); }); }); diff --git a/tests/app/managers/telegram-input-order-manager.test.ts b/tests/app/managers/telegram-input-order-manager.test.ts new file mode 100644 index 000000000..56b0403c3 --- /dev/null +++ b/tests/app/managers/telegram-input-order-manager.test.ts @@ -0,0 +1,33 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; +import { telegramInputOrderManager } from "../../../src/app/managers/telegram-input-order-manager.js"; + +describe("app/managers/telegram-input-order-manager", () => { + beforeEach(() => { + telegramInputOrderManager.__resetForTests(); + }); + + it("holds later text until every earlier album message is released", async () => { + telegramInputOrderManager.defer(777, 10); + telegramInputOrderManager.defer(777, 11); + const released = vi.fn(); + const waiting = telegramInputOrderManager.waitForEarlier(777, 12).then(released); + + await Promise.resolve(); + expect(released).not.toHaveBeenCalled(); + + telegramInputOrderManager.release(777, 10); + await Promise.resolve(); + expect(released).not.toHaveBeenCalled(); + + telegramInputOrderManager.release(777, 11); + await waiting; + expect(released).toHaveBeenCalledTimes(1); + }); + + it("does not let a deferred later update block earlier text", async () => { + telegramInputOrderManager.defer(777, 12); + + await expect(telegramInputOrderManager.waitForEarlier(777, 11)).resolves.toBeUndefined(); + }); + +}); diff --git a/tests/app/services/session-service.test.ts b/tests/app/services/session-service.test.ts index 4a5781ab4..b4db8b9a5 100644 --- a/tests/app/services/session-service.test.ts +++ b/tests/app/services/session-service.test.ts @@ -16,6 +16,11 @@ import { promptQueue } from "../../../src/app/managers/prompt-queue-manager.js"; import { promptAttachment } from "../../../src/app/managers/prompt-attachment-manager.js"; import { clearSession, setCurrentSession } from "../../../src/app/services/session-service.js"; import { createIncomingPrompt } from "../../../src/app/types/prompt.js"; +import { + clearPromptResponseMode, + consumePromptResponseMode, + setPromptResponseMode, +} from "../../../src/app/managers/prompt-response-mode-manager.js"; const SESSION = { id: "session-1", title: "Session 1", directory: "D:\\Projects\\Repo" }; @@ -24,6 +29,8 @@ describe("app/services/session-service", () => { settingsSession.current = null; promptQueue.__resetForTests(); promptAttachment.__resetForTests(); + clearPromptResponseMode("session-1"); + clearPromptResponseMode("session-2"); }); it("drops queued prompts when switching to another session", () => { @@ -35,6 +42,15 @@ describe("app/services/session-service", () => { expect(promptQueue.size()).toBe(0); }); + it("retires response modes when switching away from an abandoned session", () => { + setCurrentSession(SESSION); + setPromptResponseMode("session-1", "message-1", "text_and_tts"); + + setCurrentSession({ ...SESSION, id: "session-2" }); + + expect(consumePromptResponseMode("session-1", "message-1")).toBeNull(); + }); + it("keeps queued prompts when the same session is only renamed", () => { setCurrentSession(SESSION); promptQueue.add(createIncomingPrompt("queued for session 1")); diff --git a/tests/bot/commands/commands.test.ts b/tests/bot/commands/commands.test.ts index 5da8db279..fc2478215 100644 --- a/tests/bot/commands/commands.test.ts +++ b/tests/bot/commands/commands.test.ts @@ -41,6 +41,7 @@ const mocked = vi.hoisted(() => ({ ensureEventSubscriptionMock: vi.fn(), safeBackgroundTaskMock: vi.fn(), suppressionRegisterMock: vi.fn(), + suppressionDiscardMock: vi.fn(), attachToSessionMock: vi.fn(), })); @@ -107,7 +108,8 @@ vi.mock("../../../src/utils/safe-background-task.js", () => ({ vi.mock("../../../src/app/managers/external-input-suppression-manager.js", () => ({ externalUserInputSuppressionManager: { - register: mocked.suppressionRegisterMock, + registerMessage: mocked.suppressionRegisterMock, + discardMessage: mocked.suppressionDiscardMock, }, })); @@ -222,6 +224,7 @@ describe("bot/commands/commands", () => { mocked.ensureEventSubscriptionMock.mockReset(); mocked.safeBackgroundTaskMock.mockReset(); mocked.suppressionRegisterMock.mockReset(); + mocked.suppressionDiscardMock.mockReset(); mocked.attachToSessionMock.mockReset(); mocked.attachToSessionMock.mockResolvedValue({ busy: false, @@ -339,9 +342,13 @@ describe("bot/commands/commands", () => { }, ensureEventSubscription: mocked.ensureEventSubscriptionMock, }); - expect(mocked.suppressionRegisterMock).toHaveBeenCalledWith("session-1", "/poem"); + expect(mocked.suppressionRegisterMock).toHaveBeenCalledWith( + "session-1", + expect.any(String), + ); expect(mocked.sessionCommandMock).toHaveBeenCalledWith({ sessionID: "session-1", + messageID: expect.any(String), directory: "D:\\Projects\\Repo", command: "poem", arguments: "", @@ -380,10 +387,11 @@ describe("bot/commands/commands", () => { ); expect(mocked.suppressionRegisterMock).toHaveBeenCalledWith( "session-1", - "/poem about spring", + expect.any(String), ); expect(mocked.sessionCommandMock).toHaveBeenCalledWith({ sessionID: "session-1", + messageID: expect.any(String), directory: "D:\\Projects\\Repo", command: "poem", arguments: "about spring", diff --git a/tests/bot/handlers/media-group.test.ts b/tests/bot/handlers/media-group.test.ts index a706db436..3e57ef538 100644 --- a/tests/bot/handlers/media-group.test.ts +++ b/tests/bot/handlers/media-group.test.ts @@ -197,6 +197,11 @@ describe("bot/handlers/media-group", () => { it("queues an album as one item while the agent is busy", async () => { vi.spyOn(settingsStore, "getPromptQueueEnabled").mockReturnValue(true); + vi.spyOn(settingsStore, "getCurrentSession").mockReturnValue({ + id: "session-1", + title: "Session", + directory: "/repo", + }); foregroundSessionState.markBusy("session-1", "/repo"); const first = createPhotoContext({ messageId: 20, @@ -357,6 +362,44 @@ describe("bot/handlers/media-group", () => { ); }); + it("does not let a later album submit before an earlier album finishes", async () => { + let releaseFirst: () => void = () => {}; + const processPromptMock = vi + .fn() + .mockImplementationOnce( + () => + new Promise((resolve) => { + releaseFirst = () => resolve(true); + }), + ) + .mockResolvedValue(true); + const { deps } = createDeps({ processPrompt: processPromptMock }); + const handler = new MediaGroupAttachmentHandler(deps, { debounceMs: 1 }); + const first = createPhotoContext({ + messageId: 10, + smallFileId: "first-small", + largeFileId: "first-large", + }); + const second = createPhotoContext({ + messageId: 20, + smallFileId: "second-small", + largeFileId: "second-large", + }); + second.ctx.message!.media_group_id = "album-2"; + + await addToHandler(handler, first.ctx); + await addToHandler(handler, second.ctx); + const flushing = handler.flushAll(); + + await vi.waitFor(() => expect(processPromptMock).toHaveBeenCalledTimes(1)); + releaseFirst(); + await flushing; + + expect(processPromptMock).toHaveBeenCalledTimes(2); + expect(processPromptMock.mock.calls[0]?.[0]).toBe(first.ctx); + expect(processPromptMock.mock.calls[1]?.[0]).toBe(second.ctx); + }); + it("rejects the whole media group when any file is unsupported", async () => { const image = createDocumentContext({ messageId: 50, diff --git a/tests/bot/handlers/photo-handler.test.ts b/tests/bot/handlers/photo-handler.test.ts index 512dddda4..e95d0c7b2 100644 --- a/tests/bot/handlers/photo-handler.test.ts +++ b/tests/bot/handlers/photo-handler.test.ts @@ -67,6 +67,11 @@ describe("bot/handlers/photo-handler", () => { it("queues a photo without downloading it while the agent is busy", async () => { vi.spyOn(settingsStore, "getPromptQueueEnabled").mockReturnValue(true); + vi.spyOn(settingsStore, "getCurrentSession").mockReturnValue({ + id: "session-1", + title: "Session", + directory: "/repo", + }); foregroundSessionState.markBusy("session-1", "/repo"); const { ctx } = createPhotoContext("release screenshot"); const { deps, processPromptMock } = createDeps(); diff --git a/tests/bot/handlers/prompt.test.ts b/tests/bot/handlers/prompt.test.ts index 58d59a4c1..0eb896cee 100644 --- a/tests/bot/handlers/prompt.test.ts +++ b/tests/bot/handlers/prompt.test.ts @@ -27,6 +27,7 @@ const mocked = vi.hoisted(() => ({ sessionPromptAsyncMock: vi.fn(), sessionCreateMock: vi.fn(), suppressionRegisterMock: vi.fn(), + suppressionDiscardMock: vi.fn(), safeBackgroundTaskMock: vi.fn(), setSessionSummaryMock: vi.fn(), setBotAndChatIdMock: vi.fn(), @@ -144,7 +145,8 @@ vi.mock("../../../src/app/services/attach-service.js", () => ({ vi.mock("../../../src/app/managers/external-input-suppression-manager.js", () => ({ externalUserInputSuppressionManager: { - register: mocked.suppressionRegisterMock, + registerMessage: mocked.suppressionRegisterMock, + discardMessage: mocked.suppressionDiscardMock, }, })); @@ -252,7 +254,10 @@ describe("bot/handlers/prompt", () => { }, ensureEventSubscription: expect.any(Function), }); - expect(mocked.suppressionRegisterMock).toHaveBeenCalledWith("session-1", "Review README"); + expect(mocked.suppressionRegisterMock).toHaveBeenCalledWith( + "session-1", + expect.any(String), + ); }); it("starts prompts through promptAsync instead of the streaming prompt endpoint", async () => { @@ -273,6 +278,7 @@ describe("bot/handlers/prompt", () => { modelID: "gpt-5", }, variant: "default", + messageID: expect.stringMatching(/^msg_/), }); expect(mocked.sessionPromptMock).not.toHaveBeenCalled(); }); @@ -294,7 +300,7 @@ describe("bot/handlers/prompt", () => { ); }); - it("still notifies the user when promptAsync rejects before the run starts", async () => { + it("notifies the user while retaining delivery state after an ambiguous transport failure", async () => { const ctx = createContext(); const deps = createDeps(); @@ -314,6 +320,7 @@ describe("bot/handlers/prompt", () => { 777, "Failed to send request to OpenCode.", ); + expect(mocked.suppressionDiscardMock).not.toHaveBeenCalled(); }); it("does not notify the user when promptAsync reports an error after detach", async () => { @@ -405,7 +412,10 @@ describe("bot/handlers/prompt", () => { ]); expect(handled).toBe(true); - expect(mocked.suppressionRegisterMock).not.toHaveBeenCalled(); + expect(mocked.suppressionRegisterMock).toHaveBeenCalledWith( + "session-1", + expect.any(String), + ); }); it("keeps text prompts text-only when TTS mode is auto", async () => { @@ -414,7 +424,8 @@ describe("bot/handlers/prompt", () => { const handled = await processUserPrompt(createContext(), "Review README", createDeps()); expect(handled).toBe(true); - expect(consumePromptResponseMode("session-1")).toBe("text_only"); + const messageId = mocked.suppressionRegisterMock.mock.calls[0]?.[1] as string; + expect(consumePromptResponseMode("session-1", messageId)).toBe("text_only"); }); it("uses plural placeholder text for multiple file-only prompts", async () => { diff --git a/tests/bot/handlers/tts-response-handler.test.ts b/tests/bot/handlers/tts-response-handler.test.ts index 5e5500c31..c558f60b8 100644 --- a/tests/bot/handlers/tts-response-handler.test.ts +++ b/tests/bot/handlers/tts-response-handler.test.ts @@ -26,6 +26,7 @@ describe("bot/handlers/tts-response-handler", () => { const result = await sendTtsResponseForSession({ api: { sendAudio: sendAudioMock, sendMessage: sendMessageMock }, sessionId: "session-1", + promptMessageId: "prompt-1", chatId: 123, text: "Hello from audio", consumeResponseMode: () => "text_and_tts", @@ -51,6 +52,7 @@ describe("bot/handlers/tts-response-handler", () => { const result = await sendTtsResponseForSession({ api: { sendAudio: sendAudioMock, sendMessage: sendMessageMock }, sessionId: "session-1", + promptMessageId: "prompt-1", chatId: 123, text: "Hello from text", consumeResponseMode: () => "text_only", @@ -72,6 +74,7 @@ describe("bot/handlers/tts-response-handler", () => { const result = await sendTtsResponseForSession({ api: { sendAudio: sendAudioMock, sendMessage: sendMessageMock }, sessionId: "session-1", + promptMessageId: "prompt-1", chatId: 123, text: "Hello from audio", consumeResponseMode: () => "text_and_tts", @@ -97,6 +100,7 @@ describe("bot/handlers/tts-response-handler", () => { const result = await sendTtsResponseForSession({ api: { sendAudio: sendAudioMock, sendMessage: sendMessageMock }, sessionId: "session-1", + promptMessageId: "prompt-1", chatId: 123, text: "Hello from audio", consumeResponseMode: () => "text_and_tts", diff --git a/tests/bot/menus/inline-menu.test.ts b/tests/bot/menus/inline-menu.test.ts index e50dc9019..63990f830 100644 --- a/tests/bot/menus/inline-menu.test.ts +++ b/tests/bot/menus/inline-menu.test.ts @@ -13,6 +13,7 @@ import { defined } from "../../helpers/defined.js"; function createReplyContext(messageId: number = 1): Context { return { chat: { id: 100 }, + api: { editMessageReplyMarkup: vi.fn().mockResolvedValue(undefined) }, reply: vi.fn().mockResolvedValue({ message_id: messageId }), answerCallbackQuery: vi.fn().mockResolvedValue(undefined), deleteMessage: vi.fn().mockResolvedValue(undefined), @@ -116,6 +117,41 @@ describe("bot/menus/inline-menu", () => { }); }); + it("does not overwrite a question that wins while the menu reply is in flight", async () => { + let resolveReply: (message: { message_id: number }) => void = () => {}; + const ctx = createReplyContext(42); + (ctx.reply as ReturnType).mockImplementationOnce( + () => + new Promise<{ message_id: number }>((resolve) => { + resolveReply = resolve; + }), + ); + const keyboard = new InlineKeyboard().text("Compact output", "settings:compact_output"); + + const opening = replyWithInlineMenu(ctx, { + menuKind: "settings", + text: "Settings", + keyboard, + }); + expect(interactionManager.getSnapshot()?.metadata.reservationId).toEqual(expect.any(String)); + + interactionManager.start({ + kind: "question", + expectedInput: "text", + metadata: { requestId: "question-1" }, + }); + resolveReply({ message_id: 42 }); + + await expect(opening).resolves.toBeNull(); + expect(interactionManager.getSnapshot()).toEqual( + expect.objectContaining({ + kind: "question", + metadata: { requestId: "question-1" }, + }), + ); + expect(ctx.api.editMessageReplyMarkup).toHaveBeenCalledWith(100, 42); + }); + it("accepts callback from active inline menu", async () => { interactionManager.start({ kind: "inline", diff --git a/tests/bot/messages/external-user-input-notification.test.ts b/tests/bot/messages/external-user-input-notification.test.ts index fe5ca445a..8f96ec01e 100644 --- a/tests/bot/messages/external-user-input-notification.test.ts +++ b/tests/bot/messages/external-user-input-notification.test.ts @@ -44,6 +44,7 @@ describe("bot/messages/external-user-input-notification", () => { chatId: 777, currentSessionId: "session-1", sessionId: "session-1", + messageId: "message-1", text: "Review the parser", consumeSuppressedInput: vi.fn().mockReturnValue(false), }); @@ -66,12 +67,16 @@ describe("bot/messages/external-user-input-notification", () => { chatId: 777, currentSessionId: "session-1", sessionId: "session-1", + messageId: "message-1", text: "Review the parser", consumeSuppressedInput, }); expect(delivered).toBe(false); - expect(consumeSuppressedInput).toHaveBeenCalledWith("session-1", "Review the parser"); + expect(consumeSuppressedInput).toHaveBeenCalledWith( + "session-1", + "message-1", + ); expect(mocked.sendBotTextMock).not.toHaveBeenCalled(); }); @@ -81,6 +86,7 @@ describe("bot/messages/external-user-input-notification", () => { chatId: 777, currentSessionId: "session-2", sessionId: "session-1", + messageId: "message-1", text: "Review the parser", consumeSuppressedInput: vi.fn().mockReturnValue(false), }); diff --git a/tests/bot/middleware/interaction-guard-decision.test.ts b/tests/bot/middleware/interaction-guard-decision.test.ts index 802cdedc2..a8e252570 100644 --- a/tests/bot/middleware/interaction-guard-decision.test.ts +++ b/tests/bot/middleware/interaction-guard-decision.test.ts @@ -299,13 +299,14 @@ describe("interaction guard", () => { expect(decision.inputType).toBe("other"); }); - it("allows abort, detach, status, help, and opencode_stop while busy without interaction", () => { + it("allows utilities and settings while busy without interaction", () => { foregroundSessionState.markBusy("session-1", "D:\\Projects\\Repo"); expect(resolveInteractionGuardDecision(createContext({ text: "/abort" })).allow).toBe(true); expect(resolveInteractionGuardDecision(createContext({ text: "/detach" })).allow).toBe(true); expect(resolveInteractionGuardDecision(createContext({ text: "/status" })).allow).toBe(true); expect(resolveInteractionGuardDecision(createContext({ text: "/help" })).allow).toBe(true); + expect(resolveInteractionGuardDecision(createContext({ text: "/settings" })).allow).toBe(true); expect(resolveInteractionGuardDecision(createContext({ text: "/opencode_stop" })).allow).toBe( true, ); @@ -321,6 +322,46 @@ describe("interaction guard", () => { expect(blockedDecision.busy).toBe(true); }); + it("allows settings callbacks while busy but keeps other inline menus guarded", () => { + foregroundSessionState.markBusy("session-1", "D:\\Projects\\Repo"); + interactionManager.start({ + kind: "inline", + expectedInput: "callback", + metadata: { menuKind: "settings", messageId: 10 }, + }); + + const settingsDecision = resolveInteractionGuardDecision( + createContext({ callbackData: "settings:prompt_queue" }), + ); + expect(settingsDecision.allow).toBe(true); + expect(settingsDecision.busy).toBe(true); + + interactionManager.start({ + kind: "inline", + expectedInput: "callback", + metadata: { menuKind: "model", messageId: 11 }, + }); + const modelDecision = resolveInteractionGuardDecision( + createContext({ callbackData: "model:openai:gpt-5" }), + ); + expect(modelDecision.allow).toBe(false); + expect(modelDecision.busy).toBe(true); + }); + + it("does not let settings replace a busy question interaction", () => { + foregroundSessionState.markBusy("session-1", "D:\\Projects\\Repo"); + interactionManager.start({ + kind: "question", + expectedInput: "mixed", + }); + + const decision = resolveInteractionGuardDecision(createContext({ text: "/settings" })); + + expect(decision.allow).toBe(false); + expect(decision.reason).toBe("command_not_allowed"); + expect(decision.state?.kind).toBe("question"); + }); + it("allows opencode_stop during an active interaction without busy", () => { interactionManager.start({ kind: "inline", diff --git a/tests/bot/routers/telegram-input-order-routing.test.ts b/tests/bot/routers/telegram-input-order-routing.test.ts new file mode 100644 index 000000000..1c7559d45 --- /dev/null +++ b/tests/bot/routers/telegram-input-order-routing.test.ts @@ -0,0 +1,200 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { Composer, Context, type Api, type NextFunction } from "grammy"; +import type { Update, UserFromGetMe } from "grammy/types"; +import { telegramInputOrderManager, telegramInputOrderMiddleware } from "../../../src/app/managers/telegram-input-order-manager.js"; +import { createMediaGroupAttachmentMiddleware } from "../../../src/bot/handlers/media-group-handler.js"; +import * as sessionService from "../../../src/app/services/session-service.js"; +import * as settingsStore from "../../../src/app/stores/settings-store.js"; + +const BOT_INFO = { + id: 999, + is_bot: true, + first_name: "test", + username: "test_bot", + can_join_groups: true, + can_read_all_group_messages: false, + supports_inline_queries: false, + can_connect_to_business: false, + has_main_web_app: false, + has_topics_enabled: false, + allows_users_to_create_topics: false, + can_manage_bots: false, + supports_join_request_queries: false, +} satisfies UserFromGetMe; + +const MODEL_CAPABILITIES = { + temperature: true, + reasoning: true, + attachment: true, + toolcall: true, + input: { text: true, audio: false, image: true, video: false, pdf: true }, + output: { text: true, audio: false, image: false, video: false, pdf: false }, + interleaved: false, +}; + +function messageUpdate(updateId: number, message: Record): Update { + return { + update_id: updateId, + message: { + message_id: updateId, + date: 1, + chat: { id: 777, type: "private", first_name: "User" }, + from: { id: 1, is_bot: false, first_name: "User" }, + ...message, + }, + } as Update; +} + +describe("bot/routers Telegram album ordering", () => { + beforeEach(() => { + telegramInputOrderManager.__resetForTests(); + }); + + afterEach(() => { + vi.restoreAllMocks(); + telegramInputOrderManager.__resetForTests(); + }); + + it("submits an album under captured session A before /new and every later media route", async () => { + let currentSession = { id: "session-a", title: "A", directory: "/repo-a" }; + const currentProject = { id: "project-a", worktree: "/repo-a" }; + vi.spyOn(sessionService, "getCurrentSession").mockImplementation(() => currentSession); + vi.spyOn(settingsStore, "getCurrentProject").mockImplementation(() => currentProject); + + const order: string[] = []; + const composer = new Composer(); + composer.use(telegramInputOrderMiddleware); + composer.command("new", () => { + currentSession = { id: "session-b", title: "B", directory: "/repo-b" }; + order.push("new:session-b"); + }); + composer.on("message:voice", () => { + order.push(`voice:${currentSession.id}`); + }); + composer.on("message:audio", () => { + order.push(`audio:${currentSession.id}`); + }); + composer.on( + "message", + createMediaGroupAttachmentMiddleware( + { + bot: {} as never, + ensureEventSubscription: vi.fn(), + downloadFile: vi.fn(async () => ({ buffer: Buffer.from("image"), filePath: "image" })), + getModelCapabilities: vi.fn(async () => MODEL_CAPABILITIES), + getStoredModel: vi.fn(() => ({ providerID: "test", modelID: "test" })), + processPrompt: vi.fn(async (_ctx, _input, _deps, options) => { + order.push(`album:${options?.target?.sessionId}:${options?.target?.directory}`); + return true; + }), + }, + { debounceMs: 10 }, + ), + ); + composer.on("message:photo", () => { + order.push(`photo:${currentSession.id}`); + }); + composer.on("message:document", () => { + order.push(`document:${currentSession.id}`); + }); + + const api = { + sendMessage: vi.fn(async () => ({ message_id: 100, date: 1, chat: { id: 777, type: "private" }, text: "ok" })), + } as unknown as Api; + const dispatch = async (update: Update): Promise => { + const ctx = new Context(update, api, BOT_INFO); + await composer.middleware()(ctx, vi.fn() as unknown as NextFunction); + }; + + await dispatch( + messageUpdate(10, { + media_group_id: "album-1", + photo: [{ file_id: "photo", file_unique_id: "photo-u", width: 10, height: 10 }], + }), + ); + + await Promise.all([ + dispatch(messageUpdate(11, { text: "/new", entities: [{ type: "bot_command", offset: 0, length: 4 }] })), + dispatch(messageUpdate(12, { voice: { file_id: "voice", file_unique_id: "voice-u", duration: 1 } })), + dispatch(messageUpdate(13, { audio: { file_id: "audio", file_unique_id: "audio-u", duration: 1 } })), + dispatch(messageUpdate(14, { photo: [{ file_id: "single", file_unique_id: "single-u", width: 10, height: 10 }] })), + dispatch(messageUpdate(15, { document: { file_id: "doc", file_unique_id: "doc-u" } })), + ]); + + expect(order[0]).toBe("album:session-a:/repo-a"); + expect(order).toEqual( + expect.arrayContaining([ + "new:session-b", + "voice:session-b", + "audio:session-b", + "photo:session-b", + "document:session-b", + ]), + ); + }); + + it("allows a project-switch callback while an earlier album is pending", async () => { + let currentSession = { id: "session-a", title: "A", directory: "/repo-a" }; + let currentProject = { id: "project-a", worktree: "/repo-a" }; + vi.spyOn(sessionService, "getCurrentSession").mockImplementation(() => currentSession); + vi.spyOn(settingsStore, "getCurrentProject").mockImplementation(() => currentProject); + + const order: string[] = []; + const composer = new Composer(); + composer.use(telegramInputOrderMiddleware); + composer.callbackQuery("project:switch", () => { + currentProject = { id: "project-b", worktree: "/repo-b" }; + currentSession = { id: "session-b", title: "B", directory: "/repo-b" }; + order.push("project:session-b"); + }); + composer.on( + "message", + createMediaGroupAttachmentMiddleware( + { + bot: {} as never, + ensureEventSubscription: vi.fn(), + downloadFile: vi.fn(async () => ({ buffer: Buffer.from("image"), filePath: "image" })), + getModelCapabilities: vi.fn(async () => MODEL_CAPABILITIES), + getStoredModel: vi.fn(() => ({ providerID: "test", modelID: "test" })), + processPrompt: vi.fn(async (_ctx, _input, _deps, options) => { + order.push(`album:${options?.target?.sessionId}`); + return true; + }), + }, + { debounceMs: 10 }, + ), + ); + + const api = { + sendMessage: vi.fn(async () => ({ message_id: 100, date: 1, chat: { id: 777, type: "private" }, text: "ok" })), + } as unknown as Api; + const dispatch = async (update: Update): Promise => { + await composer.middleware()(new Context(update, api, BOT_INFO), vi.fn() as unknown as NextFunction); + }; + + await dispatch( + messageUpdate(20, { + media_group_id: "album-2", + photo: [{ file_id: "photo", file_unique_id: "photo-u", width: 10, height: 10 }], + }), + ); + const callback = dispatch({ + update_id: 21, + callback_query: { + id: "callback-1", + chat_instance: "chat-instance", + data: "project:switch", + from: { id: 1, is_bot: false, first_name: "User" }, + message: { + message_id: 19, + date: 1, + chat: { id: 777, type: "private", first_name: "User" }, + text: "Projects", + }, + }, + } as Update); + + await callback; + expect(order).toEqual(["project:session-b"]); + }); +}); diff --git a/tests/bot/services/event-subscription-service.lifecycle.test.ts b/tests/bot/services/event-subscription-service.lifecycle.test.ts index 35d4219ea..47975f930 100644 --- a/tests/bot/services/event-subscription-service.lifecycle.test.ts +++ b/tests/bot/services/event-subscription-service.lifecycle.test.ts @@ -92,7 +92,7 @@ function emitAssistantTextPart(aggregator: Aggregator, text: string): void { } as unknown as Event); } -function emitAssistantCompleted(aggregator: Aggregator): void { +function emitAssistantCompleted(aggregator: Aggregator, parentMessageId?: string): void { aggregator.processEvent({ type: "message.updated", properties: { @@ -103,6 +103,7 @@ function emitAssistantCompleted(aggregator: Aggregator): void { agent: "test-agent", providerID: "test-provider", modelID: "test-model", + ...(parentMessageId ? { parentID: parentMessageId } : {}), time: { created: Date.now() - 1000, completed: Date.now() }, }, }, @@ -840,6 +841,25 @@ describe("bot/services/event-subscription-service lifecycle", () => { ); expect(assistantRunState.finishRun("session-1", "assertion")).toBeNull(); }, 30_000); + + it("keeps parent B's TTS mode when delivery for parent A fails", async () => { + const { api, summaryAggregator } = await setupService(); + const { consumePromptResponseMode, setPromptResponseMode } = + await import("../../../src/bot/handlers/prompt.js"); + setPromptResponseMode("session-1", "parent-a", "text_only"); + setPromptResponseMode("session-1", "parent-b", "text_and_tts"); + api.sendMessage.mockRejectedValue(new Error("telegram unreachable")); + api.editMessageText.mockRejectedValue(new Error("telegram unreachable")); + + emitAssistantTextPart(summaryAggregator, "Answer A"); + emitAssistantCompleted(summaryAggregator, "parent-a"); + + await vi.waitFor(() => { + expect(api.sendMessage).toHaveBeenCalled(); + }); + expect(consumePromptResponseMode("session-1", "parent-a")).toBeNull(); + expect(consumePromptResponseMode("session-1", "parent-b")).toBe("text_and_tts"); + }); }); describe("assistant stream resilience", () => { @@ -892,7 +912,7 @@ describe("bot/services/event-subscription-service lifecycle", () => { const { api, summaryAggregator } = await setupService(); const { externalUserInputSuppressionManager } = await import("../../../src/app/managers/external-input-suppression-manager.js"); - externalUserInputSuppressionManager.register("session-1", "sent from Telegram"); + externalUserInputSuppressionManager.registerMessage("session-1", "user-message-1"); emitExternalUserMessage(summaryAggregator, "sent from Telegram"); await settle();