From e12e1387275613788f996aa59185bf05f54f0fd1 Mon Sep 17 00:00:00 2001 From: Marenz Date: Sun, 6 Sep 2026 12:26:06 +0200 Subject: [PATCH 1/5] fix(queue): retire identity suppression predictably Bind Telegram-origin suppression to message IDs instead of text so identical external input is never hidden. Expire missed-event entries on wall clock time and clear them across session and runtime teardown. --- .../external-input-suppression-manager.ts | 96 +++++++++++++++++++ .../services/external-user-input-service.ts | 6 ++ src/app/services/session-service.ts | 11 ++- src/bot/handlers/prompt.ts | 15 ++- .../external-user-input-notification.ts | 7 +- .../services/event-subscription-service.ts | 13 ++- ...external-input-suppression-manager.test.ts | 78 ++++++++++++++- tests/bot/handlers/prompt.test.ts | 23 +++-- .../external-user-input-notification.test.ts | 9 +- ...ent-subscription-service.lifecycle.test.ts | 2 +- 10 files changed, 234 insertions(+), 26 deletions(-) diff --git a/src/app/managers/external-input-suppression-manager.ts b/src/app/managers/external-input-suppression-manager.ts index 288e6ef1e..b34fcf734 100644 --- a/src/app/managers/external-input-suppression-manager.ts +++ b/src/app/managers/external-input-suppression-manager.ts @@ -1,4 +1,5 @@ const SUPPRESSION_TTL_MS = 60_000; +const MESSAGE_SUPPRESSION_TTL_MS = 5 * 60_000; interface SuppressionEntry { text: string; @@ -11,6 +12,8 @@ function normalizeExternalUserInputText(text: string): string { 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); @@ -51,14 +54,71 @@ class ExternalUserInputSuppressionManager { return true; } + registerMessage(sessionId: string, messageId: string, now: number = Date.now()): void { + if (!sessionId || !messageId) { + return; + } + + 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); + } + + consumeMessage(sessionId: string, messageId: string, _fallbackText: 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; + } + + return false; + } + + discardMessage(sessionId: string, messageId: string): void { + const messageIds = this.messageIdsBySession.get(sessionId); + if (!messageIds) { + return; + } + + messageIds.delete(messageId); + if (messageIds.size === 0) { + this.messageIdsBySession.delete(sessionId); + } + this.scheduleMessagePrune(Date.now()); + } + + clearSession(sessionId: string): void { + this.entriesBySession.delete(sessionId); + 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(); } + __getMessageCountForTests(): number { + return Array.from(this.messageIdsBySession.values()).reduce( + (count, messageIds) => count + messageIds.size, + 0, + ); + } + private prune(now: number): void { for (const [sessionId, sessionEntries] of this.entriesBySession.entries()) { const activeEntries = sessionEntries.filter((entry) => now - entry.createdAt <= SUPPRESSION_TTL_MS); @@ -70,6 +130,42 @@ class ExternalUserInputSuppressionManager { this.entriesBySession.set(sessionId, activeEntries); } } + + 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); + } + } + } + + 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?.(); + } } export const externalUserInputSuppressionManager = new ExternalUserInputSuppressionManager(); diff --git a/src/app/services/external-user-input-service.ts b/src/app/services/external-user-input-service.ts index 53f360d17..0d8b6c4f4 100644 --- a/src/app/services/external-user-input-service.ts +++ b/src/app/services/external-user-input-service.ts @@ -8,6 +8,12 @@ export interface ExternalUserInputNotification { rawFallbackText: string; } +export type ConsumeSuppressedInput = ( + sessionId: string, + messageId: string, + text: 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..d64a6b641 100644 --- a/src/app/services/session-service.ts +++ b/src/app/services/session-service.ts @@ -6,15 +6,20 @@ 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"; 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); + } } setSettingsSession(sessionInfo); @@ -25,7 +30,11 @@ 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); + } clearSettingsSession(); } diff --git a/src/bot/handlers/prompt.ts b/src/bot/handlers/prompt.ts index 03d8242a9..46af958aa 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 { randomUUID } from "node:crypto"; import { opencodeClient } from "../../opencode/client.js"; import { clearSession, @@ -370,10 +371,8 @@ export async function processUserPrompt( configuredModelID: storedModel.modelID, }); setPromptResponseMode(currentSession.id, responseMode); - - if (preparedInput.text.trim().length > 0) { - externalUserInputSuppressionManager.register(currentSession.id, preparedInput.text); - } + const promptMessageId = randomUUID(); + 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 +381,15 @@ 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); + externalUserInputSuppressionManager.discardMessage(currentSession.id, promptMessageId); const details = formatErrorDetails(error, 6000); logger.error( "[Bot] OpenCode API returned an error for session.promptAsync", @@ -407,10 +408,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/messages/external-user-input-notification.ts b/src/bot/messages/external-user-input-notification.ts index 3dfe22f5f..8f60dba10 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, text)) { return false; } diff --git a/src/bot/services/event-subscription-service.ts b/src/bot/services/event-subscription-service.ts index e2fe62caa..16c6c1ac6 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 { @@ -707,7 +708,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 +720,14 @@ 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, incomingText) => + externalUserInputSuppressionManager.consumeMessage( + incomingSessionId, + incomingMessageId, + incomingText, + ), }); } catch (err) { logger.error("[Bot] Failed to deliver external user input to Telegram:", err); @@ -1164,6 +1170,7 @@ class EventSubscriptionService implements BotEventSubscriptionService { const completedRun = assistantRunState.finishRun(sessionId, "session_idle"); clearPromptResponseMode(sessionId); + externalUserInputSuppressionManager.clearSession(sessionId); if (!this.botInstance || !this.chatIdInstance) { foregroundSessionState.markIdle(sessionId); diff --git a/tests/app/managers/external-input-suppression-manager.test.ts b/tests/app/managers/external-input-suppression-manager.test.ts index cae58110d..5864fd6a9 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,6 +6,10 @@ describe("external-input/suppression", () => { externalUserInputSuppressionManager.__resetForTests(); }); + afterEach(() => { + vi.useRealTimers(); + }); + it("consumes a matching suppressed input for the same session", () => { externalUserInputSuppressionManager.register("session-1", "Review README"); @@ -32,4 +36,76 @@ describe("external-input/suppression", () => { false, ); }); + + 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.consumeMessage( + "session-1", + "message-1", + "different rendered text", + ), + ).toBe(false); + }); + + 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.consumeMessage("session-1", "message-1", "different"), + ).toBe(true); + expect(externalUserInputSuppressionManager.__getMessageCountForTests()).toBe(0); + }); + + it("does not suppress another message with the same text", () => { + externalUserInputSuppressionManager.registerMessage("session-1", "message-1"); + + expect( + externalUserInputSuppressionManager.consumeMessage("session-1", "external-message", "same"), + ).toBe(false); + }); + + it("does not fall back to text suppression when an identity differs", () => { + externalUserInputSuppressionManager.register("session-1", "same text"); + externalUserInputSuppressionManager.registerMessage("session-1", "message-1"); + + expect( + externalUserInputSuppressionManager.consumeMessage( + "session-1", + "external-message", + "same text", + ), + ).toBe(false); + }); + + it("does not suppress the same identity in another session", () => { + externalUserInputSuppressionManager.registerMessage("session-1", "message-1"); + + expect( + externalUserInputSuppressionManager.consumeMessage("session-2", "message-1", "unrelated"), + ).toBe(false); + expect( + externalUserInputSuppressionManager.consumeMessage("session-1", "message-1", "original"), + ).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.__getMessageCountForTests()).toBe(0); + }); }); diff --git a/tests/bot/handlers/prompt.test.ts b/tests/bot/handlers/prompt.test.ts index 58d59a4c1..c6a3776e4 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.any(String), }); 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("retains delivery state when promptAsync has an ambiguous transport failure", async () => { const ctx = createContext(); const deps = createDeps(); @@ -310,10 +316,8 @@ describe("bot/handlers/prompt", () => { backgroundTask.onError?.(error); }); - expect(deps.bot.api.sendMessage).toHaveBeenCalledWith( - 777, - "Failed to send request to OpenCode.", - ); + expect(deps.bot.api.sendMessage).not.toHaveBeenCalled(); + expect(mocked.suppressionDiscardMock).not.toHaveBeenCalled(); }); it("does not notify the user when promptAsync reports an error after detach", async () => { @@ -405,7 +409,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 () => { diff --git a/tests/bot/messages/external-user-input-notification.test.ts b/tests/bot/messages/external-user-input-notification.test.ts index fe5ca445a..de5facabc 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,17 @@ 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", + "Review the parser", + ); expect(mocked.sendBotTextMock).not.toHaveBeenCalled(); }); @@ -81,6 +87,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/services/event-subscription-service.lifecycle.test.ts b/tests/bot/services/event-subscription-service.lifecycle.test.ts index 35d4219ea..9da1674ad 100644 --- a/tests/bot/services/event-subscription-service.lifecycle.test.ts +++ b/tests/bot/services/event-subscription-service.lifecycle.test.ts @@ -892,7 +892,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(); From 77203e43b6e846ab2a7875d969e22060b2759de0 Mon Sep 17 00:00:00 2001 From: Marenz Date: Sun, 6 Sep 2026 12:28:06 +0200 Subject: [PATCH 2/5] fix(queue): scope response modes to parent prompts Track text and TTS delivery by session and parent message identity. A failed response can now retire only its own mode without stripping TTS from another queued response. --- .../managers/prompt-response-mode-manager.ts | 35 +++++++++++++++++++ .../managers/summary-aggregation-manager.ts | 2 ++ src/app/services/session-service.ts | 3 ++ src/app/services/tts-service.ts | 6 ++-- src/bot/handlers/prompt.ts | 35 +++++++++---------- src/bot/handlers/tts-response-handler.ts | 12 ++++++- .../services/event-subscription-service.ts | 18 ++++++---- tests/app/services/session-service.test.ts | 16 +++++++++ tests/bot/handlers/prompt.test.ts | 3 +- .../bot/handlers/tts-response-handler.test.ts | 4 +++ ...ent-subscription-service.lifecycle.test.ts | 22 +++++++++++- 11 files changed, 126 insertions(+), 30 deletions(-) create mode 100644 src/app/managers/prompt-response-mode-manager.ts 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/services/session-service.ts b/src/app/services/session-service.ts index d64a6b641..b1a13e16a 100644 --- a/src/app/services/session-service.ts +++ b/src/app/services/session-service.ts @@ -7,6 +7,7 @@ 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 }; @@ -19,6 +20,7 @@ export function setCurrentSession(sessionInfo: SessionInfo): void { promptAttachment.clear("session_switched"); if (previousSessionId) { externalUserInputSuppressionManager.clearSession(previousSessionId); + clearPromptResponseMode(previousSessionId); } } @@ -35,6 +37,7 @@ export function clearSession(): void { 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/handlers/prompt.ts b/src/bot/handlers/prompt.ts index 46af958aa..f5da1cc6a 100644 --- a/src/bot/handlers/prompt.ts +++ b/src/bot/handlers/prompt.ts @@ -45,13 +45,22 @@ 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"; type ProcessPromptOptions = { responseMode?: PromptResponseMode; @@ -65,18 +74,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 { @@ -370,8 +370,8 @@ export async function processUserPrompt( configuredProviderID: storedModel.providerID, configuredModelID: storedModel.modelID, }); - setPromptResponseMode(currentSession.id, responseMode); const promptMessageId = randomUUID(); + setPromptResponseMode(currentSession.id, promptMessageId, responseMode); externalUserInputSuppressionManager.registerMessage(currentSession.id, promptMessageId); // CRITICAL: Use the async prompt start endpoint here. @@ -388,8 +388,7 @@ export async function processUserPrompt( foregroundSessionState.markIdle(currentSession.id); void markAttachedSessionIdle(currentSession.id); assistantRunState.clearRun(currentSession.id, "session_prompt_api_error"); - clearPromptResponseMode(currentSession.id); - externalUserInputSuppressionManager.discardMessage(currentSession.id, promptMessageId); + discardPromptDeliveryState(currentSession.id, promptMessageId); const details = formatErrorDetails(error, 6000); logger.error( "[Bot] OpenCode API returned an error for session.promptAsync", 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/services/event-subscription-service.ts b/src/bot/services/event-subscription-service.ts index 16c6c1ac6..f9b1ea2f8 100644 --- a/src/bot/services/event-subscription-service.ts +++ b/src/bot/services/event-subscription-service.ts @@ -608,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"); @@ -621,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"); @@ -690,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"); @@ -1225,7 +1232,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); @@ -1234,8 +1240,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"); @@ -1246,7 +1251,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/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/handlers/prompt.test.ts b/tests/bot/handlers/prompt.test.ts index c6a3776e4..70876b267 100644 --- a/tests/bot/handlers/prompt.test.ts +++ b/tests/bot/handlers/prompt.test.ts @@ -421,7 +421,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/services/event-subscription-service.lifecycle.test.ts b/tests/bot/services/event-subscription-service.lifecycle.test.ts index 9da1674ad..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", () => { From a435ad8112eb4658633795f41de6bfac6a13529a Mon Sep 17 00:00:00 2001 From: Marenz Date: Sun, 6 Sep 2026 12:30:29 +0200 Subject: [PATCH 3/5] fix(queue): preserve Telegram album input order Defer later commands, callbacks, voice, audio, photos, documents, and albums behind earlier album batches. Capture the original project and session so a delayed album cannot land in a newly selected context. --- src/app/managers/interaction-manager.ts | 26 +++ .../managers/telegram-input-order-manager.ts | 108 ++++++++++ .../command-catalog-callback-handler.ts | 9 +- src/bot/handlers/media-group-handler.ts | 42 +++- src/bot/index.ts | 2 + src/bot/middleware/interaction-guard.ts | 5 + .../telegram-input-order-manager.test.ts | 45 ++++ tests/bot/commands/commands.test.ts | 14 +- tests/bot/handlers/media-group.test.ts | 43 ++++ tests/bot/handlers/photo-handler.test.ts | 5 + .../telegram-input-order-routing.test.ts | 200 ++++++++++++++++++ 11 files changed, 490 insertions(+), 9 deletions(-) create mode 100644 src/app/managers/telegram-input-order-manager.ts create mode 100644 tests/app/managers/telegram-input-order-manager.test.ts create mode 100644 tests/bot/routers/telegram-input-order-routing.test.ts 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/telegram-input-order-manager.ts b/src/app/managers/telegram-input-order-manager.ts new file mode 100644 index 000000000..c16ef31b1 --- /dev/null +++ b/src/app/managers/telegram-input-order-manager.ts @@ -0,0 +1,108 @@ +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); + }); + } + + waitForPending(chatId: number): Promise { + return this.waitForEarlier(chatId, Number.POSITIVE_INFINITY); + } + + __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 telegramInputOrderManager.waitForPending(chatId); + await next(); + return; + } + + if (message.media_group_id) { + await next(); + return; + } + + await telegramInputOrderManager.waitForEarlier(chatId, message.message_id); + await next(); +} diff --git a/src/bot/callbacks/command-catalog-callback-handler.ts b/src/bot/callbacks/command-catalog-callback-handler.ts index 2dbdf7b31..a5cbe4747 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 { randomUUID } from "node:crypto"; 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 = randomUUID(); + 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/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/middleware/interaction-guard.ts b/src/bot/middleware/interaction-guard.ts index d2286fb6d..f1781f42c 100644 --- a/src/bot/middleware/interaction-guard.ts +++ b/src/bot/middleware/interaction-guard.ts @@ -5,6 +5,7 @@ import { reconcileForegroundBusyState } from "../../app/services/run-control-ser import { canQueueMediaPrompt, rejectQueuedMediaBeforePreparation, + dispatchNextQueuedPrompt, shouldSuggestPromptQueue, tryEnqueuePrompt, } from "../handlers/prompt-queue-dispatch.js"; @@ -12,6 +13,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 +145,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/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..436c5027a --- /dev/null +++ b/tests/app/managers/telegram-input-order-manager.test.ts @@ -0,0 +1,45 @@ +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(); + }); + + it("holds a callback without a message id until the album is released", async () => { + telegramInputOrderManager.defer(777, 10); + const released = vi.fn(); + const waiting = telegramInputOrderManager.waitForPending(777).then(released); + + await Promise.resolve(); + expect(released).not.toHaveBeenCalled(); + + telegramInputOrderManager.release(777, 10); + await waiting; + expect(released).toHaveBeenCalledTimes(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/routers/telegram-input-order-routing.test.ts b/tests/bot/routers/telegram-input-order-routing.test.ts new file mode 100644 index 000000000..24364569a --- /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("holds a project-switch callback behind an earlier album", 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(["album:session-a", "project:session-b"]); + }); +}); From 05590db7fe888c4549ff9aa9321f4905c4f25326 Mon Sep 17 00:00:00 2001 From: Marenz Date: Sun, 6 Sep 2026 12:31:34 +0200 Subject: [PATCH 4/5] fix(queue): reserve settings interactions while busy Allow /settings only when no unrelated interaction owns input, while preserving the active settings flow. Update queue wording and documentation for durable OpenCode admission semantics. --- PRODUCT.md | 2 +- src/bot/handlers/prompt-queue-dispatch.ts | 5 +- src/bot/menus/inline-menu.ts | 36 +++++++++++++-- .../middleware/interaction-guard-decision.ts | 46 +++++++++++++++---- src/i18n/ar.ts | 6 +-- src/i18n/de.ts | 7 ++- src/i18n/en.ts | 7 +-- src/i18n/es.ts | 7 ++- src/i18n/fr.ts | 7 ++- src/i18n/it.ts | 6 +-- src/i18n/ko.ts | 6 +-- src/i18n/pt.ts | 7 ++- src/i18n/ru.ts | 6 +-- src/i18n/zh.ts | 6 +-- tests/bot/menus/inline-menu.test.ts | 36 +++++++++++++++ .../interaction-guard-decision.test.ts | 43 ++++++++++++++++- .../compact-progress-streamer.test.ts | 8 ++-- 17 files changed, 184 insertions(+), 57 deletions(-) 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/bot/handlers/prompt-queue-dispatch.ts b/src/bot/handlers/prompt-queue-dispatch.ts index 3d7b4f01f..474a88f5f 100644 --- a/src/bot/handlers/prompt-queue-dispatch.ts +++ b/src/bot/handlers/prompt-queue-dispatch.ts @@ -94,10 +94,7 @@ export async function tryEnqueuePrompt(ctx: Context, input: QueuedPromptInput): logger.info( `[PromptQueue] Prompt queued while session is busy: size=${promptQueue.size()}/${MAX_QUEUED_PROMPTS}`, ); - await replyWithKeyboard( - ctx, - t("queue.added", { count: String(promptQueue.size()), max: String(MAX_QUEUED_PROMPTS) }), - ); + await replyWithKeyboard(ctx, t("queue.added")); return true; } diff --git a/src/bot/menus/inline-menu.ts b/src/bot/menus/inline-menu.ts index f09144f88..653387864 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,23 @@ export function appendInlineMenuCancelButton( export async function replyWithInlineMenu( ctx: Context, options: InlineMenuReplyOptions, -): Promise { +): Promise { const keyboard = appendInlineMenuCancelButton(options.keyboard, options.menuKind); + 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 +130,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/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/i18n/ar.ts b/src/i18n/ar.ts index 4228364a7..748c77e18 100644 --- a/src/i18n/ar.ts +++ b/src/i18n/ar.ts @@ -68,7 +68,7 @@ export const ar: I18nDictionary = { "bot.thinking": "💭 جارٍ التفكير...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ جارٍ العمل", - "progress.compact.finished_header": "✅ اكتمل العمل", + "progress.compact.finished_header": "✅ اكتمل الرد", "progress.compact.thinking": "💭 جارٍ التفكير...", "progress.compact.responding": "✍️ جارٍ كتابة الرد...", "progress.compact.waiting_question": "❓ في انتظار إجابتك...", @@ -391,8 +391,8 @@ export const ar: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 الإعداد الافتراضي", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 أُضيفت إلى قائمة الانتظار ({count}/{max}). ستُرسل بعد انتهاء المهمة الحالية.", - "queue.full": "⚠️ قائمة الانتظار ممتلئة ({max}). احذف رسالة أو انتظر انتهاء المهمة الحالية.", + "queue.added": "📥 جارٍ الإرسال إلى قائمة انتظار OpenCode.", + "queue.full": "⚠️ قائمة الإرسال المعلّق ممتلئة ({max}). احذف رسالة أو حاول مرة أخرى بعد قليل.", "queue.media_limit": "⚠️ الوسائط في قائمة الانتظار محدودة بـ {maxSizeMb} MiB. انتظر إرسال عنصر ثم أعد المحاولة.", "queue.removed": "🗑 تمت إزالة الرسالة من قائمة الانتظار.", "queue.not_found": "لم تعد هذه الرسالة في قائمة الانتظار.", diff --git a/src/i18n/de.ts b/src/i18n/de.ts index 62a6d9845..f53b65308 100644 --- a/src/i18n/de.ts +++ b/src/i18n/de.ts @@ -67,7 +67,7 @@ export const de: I18nDictionary = { "bot.thinking": "💭 Denke...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Arbeite", - "progress.compact.finished_header": "✅ Arbeit abgeschlossen", + "progress.compact.finished_header": "✅ Antwort abgeschlossen", "progress.compact.thinking": "💭 Denke...", "progress.compact.responding": "✍️ Schreibe Antwort...", "progress.compact.waiting_question": "❓ Warte auf deine Antwort...", @@ -420,11 +420,10 @@ export const de: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Standard", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": - "📥 Zur Warteschlange hinzugefügt ({count}/{max}). Die Nachricht wird gesendet, sobald die aktuelle Aufgabe abgeschlossen ist.", + "queue.added": "📥 Wird an OpenCodes Warteschlange übermittelt.", "queue.media_limit": "⚠️ Medien in der Warteschlange sind auf {maxSizeMb} MiB begrenzt. Warte, bis ein Eintrag gesendet wurde.", "queue.full": - "⚠️ Die Warteschlange ist voll ({max}). Entferne eine Nachricht oder warte, bis die aktuelle Aufgabe abgeschlossen ist.", + "⚠️ Die Warteschlange ausstehender Übermittlungen ist voll ({max}). Entferne eine Nachricht oder versuche es gleich erneut.", "queue.removed": "🗑 Nachricht aus der Warteschlange entfernt.", "queue.not_found": "Diese Nachricht ist nicht mehr in der Warteschlange.", "queue.disabled_hint": "Die Nachrichtenwarteschlange lässt sich in /settings aktivieren.", diff --git a/src/i18n/en.ts b/src/i18n/en.ts index 6d4bae7ef..d15de4105 100644 --- a/src/i18n/en.ts +++ b/src/i18n/en.ts @@ -64,7 +64,7 @@ export const en = { "bot.thinking": "💭 Thinking...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Working", - "progress.compact.finished_header": "✅ Finished Work", + "progress.compact.finished_header": "✅ Response complete", "progress.compact.thinking": "💭 Thinking...", "progress.compact.responding": "✍️ Writing answer...", "progress.compact.waiting_question": "❓ Waiting for your answer...", @@ -401,8 +401,9 @@ export const en = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Default", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 Added to queue ({count}/{max}). It will be sent when the current task finishes.", - "queue.full": "⚠️ Queue is full ({max}). Remove a message or wait for the current task to finish.", + "queue.added": "📥 Forwarding to OpenCode's queue.", + "queue.full": + "⚠️ The pending submission queue is full ({max}). Remove a message or try again shortly.", "queue.media_limit": "⚠️ Queued media is limited to {maxSizeMb} MiB. Wait for an item to send, then try again.", "queue.removed": "🗑 Message removed from the queue.", "queue.not_found": "This message is no longer in the queue.", diff --git a/src/i18n/es.ts b/src/i18n/es.ts index 4ba8c672b..eb3d5a32a 100644 --- a/src/i18n/es.ts +++ b/src/i18n/es.ts @@ -67,7 +67,7 @@ export const es: I18nDictionary = { "bot.thinking": "💭 Pensando...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Trabajando", - "progress.compact.finished_header": "✅ Trabajo terminado", + "progress.compact.finished_header": "✅ Respuesta completada", "progress.compact.thinking": "💭 Pensando...", "progress.compact.responding": "✍️ Escribiendo respuesta...", "progress.compact.waiting_question": "❓ Esperando tu respuesta...", @@ -417,11 +417,10 @@ export const es: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Predeterminado", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": - "📥 Añadido a la cola ({count}/{max}). Se enviará cuando termine la tarea actual.", + "queue.added": "📥 Enviando a la cola de OpenCode.", "queue.media_limit": "⚠️ Los archivos multimedia en cola están limitados a {maxSizeMb} MiB. Espera a que se envíe un elemento.", "queue.full": - "⚠️ La cola está llena ({max}). Elimina un mensaje o espera a que termine la tarea actual.", + "⚠️ La cola de envíos pendientes está llena ({max}). Elimina un mensaje o inténtalo de nuevo en breve.", "queue.removed": "🗑 Mensaje eliminado de la cola.", "queue.not_found": "Este mensaje ya no está en la cola.", "queue.disabled_hint": "La cola de mensajes se activa en /settings.", diff --git a/src/i18n/fr.ts b/src/i18n/fr.ts index 2c5a9bc57..277c0074c 100644 --- a/src/i18n/fr.ts +++ b/src/i18n/fr.ts @@ -67,7 +67,7 @@ export const fr: I18nDictionary = { "bot.thinking": "💭 Réflexion en cours...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Travail en cours", - "progress.compact.finished_header": "✅ Travail terminé", + "progress.compact.finished_header": "✅ Réponse terminée", "progress.compact.thinking": "💭 Réflexion en cours...", "progress.compact.responding": "✍️ Rédaction de la réponse...", "progress.compact.waiting_question": "❓ En attente de votre réponse...", @@ -421,11 +421,10 @@ export const fr: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Par défaut", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": - "📥 Ajouté à la file d'attente ({count}/{max}). Le message sera envoyé à la fin de la tâche en cours.", + "queue.added": "📥 Envoi vers la file d'attente d'OpenCode.", "queue.media_limit": "⚠️ Les médias en file sont limités à {maxSizeMb} MiB. Attendez l'envoi d'un élément.", "queue.full": - "⚠️ La file d'attente est pleine ({max}). Supprimez un message ou attendez la fin de la tâche en cours.", + "⚠️ La file des envois en attente est pleine ({max}). Supprimez un message ou réessayez bientôt.", "queue.removed": "🗑 Message retiré de la file d'attente.", "queue.not_found": "Ce message n'est plus dans la file d'attente.", "queue.disabled_hint": "La file d'attente des messages s'active dans /settings.", diff --git a/src/i18n/it.ts b/src/i18n/it.ts index 53775c9d9..d2774767c 100644 --- a/src/i18n/it.ts +++ b/src/i18n/it.ts @@ -68,7 +68,7 @@ export const it: I18nDictionary = { "bot.thinking": "💭 Sto pensando...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Al lavoro", - "progress.compact.finished_header": "✅ Lavoro completato", + "progress.compact.finished_header": "✅ Risposta completata", "progress.compact.thinking": "💭 Sto pensando...", "progress.compact.responding": "✍️ Scrivo la risposta...", "progress.compact.waiting_question": "❓ In attesa della tua risposta...", @@ -416,8 +416,8 @@ export const it: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Predefinito", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 Aggiunto alla coda ({count}/{max}). Verrà inviato quando l'attività corrente termina.", - "queue.full": "⚠️ La coda è piena ({max}). Rimuovi un messaggio o attendi che l'attività corrente termini.", + "queue.added": "📥 Invio alla coda di OpenCode.", + "queue.full": "⚠️ La coda degli invii in sospeso è piena ({max}). Rimuovi un messaggio o riprova tra poco.", "queue.media_limit": "⚠️ I media in coda sono limitati a {maxSizeMb} MiB. Attendi l'invio di un elemento e riprova.", "queue.removed": "🗑 Messaggio rimosso dalla coda.", "queue.not_found": "Questo messaggio non è più in coda.", diff --git a/src/i18n/ko.ts b/src/i18n/ko.ts index 7cbb9a821..c09139b10 100644 --- a/src/i18n/ko.ts +++ b/src/i18n/ko.ts @@ -73,7 +73,7 @@ export const ko: I18nDictionary = { "bot.thinking": "💭 생각하는 중...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ 작업 중", - "progress.compact.finished_header": "✅ 작업 완료", + "progress.compact.finished_header": "✅ 응답 완료", "progress.compact.thinking": "💭 생각하는 중...", "progress.compact.responding": "✍️ 답변 작성 중...", "progress.compact.waiting_question": "❓ 답변을 기다리는 중...", @@ -410,8 +410,8 @@ export const ko: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 기본값", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 대기열에 추가되었습니다 ({count}/{max}). 현재 작업이 끝나면 전송됩니다.", - "queue.full": "⚠️ 대기열이 가득 찼습니다 ({max}). 메시지를 삭제하거나 현재 작업이 끝날 때까지 기다려 주세요.", + "queue.added": "📥 OpenCode 대기열로 전달 중입니다.", + "queue.full": "⚠️ 제출 대기열이 가득 찼습니다 ({max}). 메시지를 삭제하거나 잠시 후 다시 시도해 주세요.", "queue.media_limit": "⚠️ 대기열 미디어는 총 {maxSizeMb} MiB로 제한됩니다. 항목이 전송된 후 다시 시도하세요.", "queue.removed": "🗑 대기열에서 메시지를 삭제했습니다.", "queue.not_found": "이 메시지는 더 이상 대기열에 없습니다.", diff --git a/src/i18n/pt.ts b/src/i18n/pt.ts index 5477ce8a3..0a459213a 100644 --- a/src/i18n/pt.ts +++ b/src/i18n/pt.ts @@ -66,7 +66,7 @@ export const pt: I18nDictionary = { "bot.thinking": "💭 Pensando...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Trabalhando", - "progress.compact.finished_header": "✅ Trabalho concluído", + "progress.compact.finished_header": "✅ Resposta concluída", "progress.compact.thinking": "💭 Pensando...", "progress.compact.responding": "✍️ Escrevendo resposta...", "progress.compact.waiting_question": "❓ Aguardando sua resposta...", @@ -418,11 +418,10 @@ export const pt: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Padrão", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": - "📥 Adicionado à fila ({count}/{max}). Será enviado quando a tarefa atual terminar.", + "queue.added": "📥 Enviando para a fila do OpenCode.", "queue.media_limit": "⚠️ A mídia na fila está limitada a {maxSizeMb} MiB. Aguarde o envio de um item.", "queue.full": - "⚠️ A fila está cheia ({max}). Remova uma mensagem ou aguarde o término da tarefa atual.", + "⚠️ A fila de envios pendentes está cheia ({max}). Remova uma mensagem ou tente novamente em breve.", "queue.removed": "🗑 Mensagem removida da fila.", "queue.not_found": "Esta mensagem não está mais na fila.", "queue.disabled_hint": "A fila de mensagens pode ser ativada em /settings.", diff --git a/src/i18n/ru.ts b/src/i18n/ru.ts index f7ce7fd8d..9e6321301 100644 --- a/src/i18n/ru.ts +++ b/src/i18n/ru.ts @@ -64,7 +64,7 @@ export const ru: I18nDictionary = { "bot.thinking": "💭 Думаю...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Работаю", - "progress.compact.finished_header": "✅ Работа завершена", + "progress.compact.finished_header": "✅ Ответ завершён", "progress.compact.thinking": "💭 Думаю...", "progress.compact.responding": "✍️ Пишу ответ...", "progress.compact.waiting_question": "❓ Жду ваш ответ...", @@ -404,8 +404,8 @@ export const ru: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Default", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 Добавлено в очередь ({count}/{max}). Сообщение уйдёт после завершения текущей задачи.", - "queue.full": "⚠️ Очередь заполнена ({max}). Удалите сообщение или дождитесь завершения текущей задачи.", + "queue.added": "📥 Отправляется в очередь OpenCode.", + "queue.full": "⚠️ Очередь ожидающих отправок заполнена ({max}). Удалите сообщение или повторите попытку чуть позже.", "queue.media_limit": "⚠️ Медиа в очереди ограничены {maxSizeMb} MiB. Дождитесь отправки элемента и повторите попытку.", "queue.removed": "🗑 Сообщение удалено из очереди.", "queue.not_found": "Этого сообщения больше нет в очереди.", diff --git a/src/i18n/zh.ts b/src/i18n/zh.ts index 7678bba03..e4fa8420f 100644 --- a/src/i18n/zh.ts +++ b/src/i18n/zh.ts @@ -59,7 +59,7 @@ export const zh: I18nDictionary = { "bot.thinking": "💭 思考中...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ 工作中", - "progress.compact.finished_header": "✅ 工作完成", + "progress.compact.finished_header": "✅ 回复完成", "progress.compact.thinking": "💭 思考中...", "progress.compact.responding": "✍️ 正在撰写回复...", "progress.compact.waiting_question": "❓ 等待你的回答...", @@ -368,8 +368,8 @@ export const zh: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 默认", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 已加入队列({count}/{max})。当前任务完成后将自动发送。", - "queue.full": "⚠️ 队列已满({max})。请删除一条消息或等待当前任务完成。", + "queue.added": "📥 正在提交到 OpenCode 队列。", + "queue.full": "⚠️ 待提交队列已满({max})。请删除一条消息或稍后重试。", "queue.media_limit": "⚠️ 队列媒体总大小限制为 {maxSizeMb} MiB。请等待一个项目发送后重试。", "queue.removed": "🗑 消息已从队列中移除。", "queue.not_found": "该消息已不在队列中。", 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/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/streaming/compact-progress-streamer.test.ts b/tests/bot/streaming/compact-progress-streamer.test.ts index 6d9dfa920..10bbe6581 100644 --- a/tests/bot/streaming/compact-progress-streamer.test.ts +++ b/tests/bot/streaming/compact-progress-streamer.test.ts @@ -21,7 +21,7 @@ describe("bot/streaming/compact-progress-streamer", () => { expect(editText).toHaveBeenCalledWith( "s1", 10, - "✅ Finished Work\ntool calls: 0 · changed files: 0", + "✅ Response complete\ntool calls: 0 · changed files: 0", ); }); @@ -67,7 +67,7 @@ describe("bot/streaming/compact-progress-streamer", () => { expect(editText).toHaveBeenCalledWith( "s1", 10, - "✅ Finished Work\ntool calls: 0 · changed files: 0", + "✅ Response complete\ntool calls: 0 · changed files: 0", ); }); @@ -104,7 +104,7 @@ describe("bot/streaming/compact-progress-streamer", () => { 2, "s1", 10, - "✅ Finished Work\ntool calls: 1 · changed files: 0", + "✅ Response complete\ntool calls: 1 · changed files: 0", ); }); @@ -126,7 +126,7 @@ describe("bot/streaming/compact-progress-streamer", () => { expect(sendText).toHaveBeenCalledTimes(1); expect(sendText).toHaveBeenCalledWith( "s1", - "✅ Finished Work\ntool calls: 2 · changed files: 2", + "✅ Response complete\ntool calls: 2 · changed files: 2", ); expect(editText).not.toHaveBeenCalled(); }); From 7b2bad22aa6ec9eca59cf5efa51bf5c1a655a625 Mon Sep 17 00:00:00 2001 From: "Mathias L. Baumann" Date: Mon, 7 Sep 2026 14:52:32 +0200 Subject: [PATCH 5/5] fix(queue): restore bot-owned delivery semantics Signed-off-by: Mathias L. Baumann --- .../external-input-suppression-manager.ts | 66 +------------------ .../managers/telegram-input-order-manager.ts | 5 -- .../services/external-user-input-service.ts | 1 - .../command-catalog-callback-handler.ts | 4 +- src/bot/handlers/prompt-queue-dispatch.ts | 5 +- src/bot/handlers/prompt.ts | 21 +++++- src/bot/menus/inline-menu.ts | 4 ++ .../external-user-input-notification.ts | 2 +- src/bot/middleware/interaction-guard.ts | 1 - .../services/event-subscription-service.ts | 5 +- src/i18n/ar.ts | 6 +- src/i18n/de.ts | 7 +- src/i18n/en.ts | 7 +- src/i18n/es.ts | 7 +- src/i18n/fr.ts | 7 +- src/i18n/it.ts | 6 +- src/i18n/ko.ts | 6 +- src/i18n/pt.ts | 7 +- src/i18n/ru.ts | 6 +- src/i18n/zh.ts | 6 +- src/utils/opencode-message-id.ts | 5 ++ ...external-input-suppression-manager.test.ts | 54 ++------------- .../telegram-input-order-manager.test.ts | 12 ---- tests/bot/handlers/prompt.test.ts | 9 ++- .../external-user-input-notification.test.ts | 1 - .../telegram-input-order-routing.test.ts | 4 +- .../compact-progress-streamer.test.ts | 8 +-- 27 files changed, 88 insertions(+), 184 deletions(-) create mode 100644 src/utils/opencode-message-id.ts diff --git a/src/app/managers/external-input-suppression-manager.ts b/src/app/managers/external-input-suppression-manager.ts index b34fcf734..98e6d3508 100644 --- a/src/app/managers/external-input-suppression-manager.ts +++ b/src/app/managers/external-input-suppression-manager.ts @@ -1,59 +1,9 @@ -const SUPPRESSION_TTL_MS = 60_000; const MESSAGE_SUPPRESSION_TTL_MS = 5 * 60_000; -interface SuppressionEntry { - text: string; - createdAt: number; -} - -function normalizeExternalUserInputText(text: string): string { - return text.replace(/\r\n/g, "\n").trim(); -} - 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) { - return; - } - - this.prune(now); - - const sessionEntries = this.entriesBySession.get(sessionId) ?? []; - sessionEntries.push({ text: normalizedText, createdAt: now }); - this.entriesBySession.set(sessionId, sessionEntries); - } - - consume(sessionId: string, text: string, now: number = Date.now()): boolean { - const normalizedText = normalizeExternalUserInputText(text); - if (!sessionId || !normalizedText) { - return false; - } - - this.prune(now); - - const sessionEntries = this.entriesBySession.get(sessionId); - if (!sessionEntries?.length) { - return false; - } - - const entryIndex = sessionEntries.findIndex((entry) => entry.text === normalizedText); - if (entryIndex < 0) { - return false; - } - - sessionEntries.splice(entryIndex, 1); - if (sessionEntries.length === 0) { - this.entriesBySession.delete(sessionId); - } - - return true; - } - registerMessage(sessionId: string, messageId: string, now: number = Date.now()): void { if (!sessionId || !messageId) { return; @@ -66,7 +16,7 @@ class ExternalUserInputSuppressionManager { this.scheduleMessagePrune(now); } - consumeMessage(sessionId: string, messageId: string, _fallbackText: string): boolean { + consumeMessage(sessionId: string, messageId: string): boolean { this.pruneMessages(Date.now()); const messageIds = this.messageIdsBySession.get(sessionId); if (messageIds?.delete(messageId)) { @@ -94,13 +44,11 @@ class ExternalUserInputSuppressionManager { } clearSession(sessionId: string): void { - this.entriesBySession.delete(sessionId); this.messageIdsBySession.delete(sessionId); this.scheduleMessagePrune(Date.now()); } clearAll(): void { - this.entriesBySession.clear(); this.messageIdsBySession.clear(); if (this.messagePruneTimer) { clearTimeout(this.messagePruneTimer); @@ -119,18 +67,6 @@ class ExternalUserInputSuppressionManager { ); } - 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; - } - - this.entriesBySession.set(sessionId, activeEntries); - } - } - private pruneMessages(now: number): void { for (const [sessionId, messageIds] of this.messageIdsBySession.entries()) { for (const [messageId, expiresAt] of messageIds.entries()) { diff --git a/src/app/managers/telegram-input-order-manager.ts b/src/app/managers/telegram-input-order-manager.ts index c16ef31b1..041f9e04d 100644 --- a/src/app/managers/telegram-input-order-manager.ts +++ b/src/app/managers/telegram-input-order-manager.ts @@ -36,10 +36,6 @@ class TelegramInputOrderManager { }); } - waitForPending(chatId: number): Promise { - return this.waitForEarlier(chatId, Number.POSITIVE_INFINITY); - } - __resetForTests(): void { for (const waiters of this.waitersByChat.values()) { for (const waiter of waiters) { @@ -93,7 +89,6 @@ export async function telegramInputOrderMiddleware( } if (!message) { - await telegramInputOrderManager.waitForPending(chatId); await next(); return; } diff --git a/src/app/services/external-user-input-service.ts b/src/app/services/external-user-input-service.ts index 0d8b6c4f4..2ab9cfe71 100644 --- a/src/app/services/external-user-input-service.ts +++ b/src/app/services/external-user-input-service.ts @@ -11,7 +11,6 @@ export interface ExternalUserInputNotification { export type ConsumeSuppressedInput = ( sessionId: string, messageId: string, - text: string, ) => boolean; function normalizeExternalUserInputText(text: string): string { diff --git a/src/bot/callbacks/command-catalog-callback-handler.ts b/src/bot/callbacks/command-catalog-callback-handler.ts index a5cbe4747..1683f8a55 100644 --- a/src/bot/callbacks/command-catalog-callback-handler.ts +++ b/src/bot/callbacks/command-catalog-callback-handler.ts @@ -27,7 +27,7 @@ import { markAttachedSessionIdle, } from "../../app/services/attach-service.js"; import { externalUserInputSuppressionManager } from "../../app/managers/external-input-suppression-manager.js"; -import { randomUUID } from "node:crypto"; +import { createOpencodeMessageId } from "../../utils/opencode-message-id.js"; import { opencodeClient } from "../../opencode/client.js"; import { buildCommandsConfirmKeyboard, @@ -286,7 +286,7 @@ export async function executeCommand( configuredProviderID: storedModel.providerID, configuredModelID: storedModel.modelID, }); - const promptMessageId = randomUUID(); + const promptMessageId = createOpencodeMessageId(); externalUserInputSuppressionManager.registerMessage(session.id, promptMessageId); safeBackgroundTask({ diff --git a/src/bot/handlers/prompt-queue-dispatch.ts b/src/bot/handlers/prompt-queue-dispatch.ts index 474a88f5f..3d7b4f01f 100644 --- a/src/bot/handlers/prompt-queue-dispatch.ts +++ b/src/bot/handlers/prompt-queue-dispatch.ts @@ -94,7 +94,10 @@ export async function tryEnqueuePrompt(ctx: Context, input: QueuedPromptInput): logger.info( `[PromptQueue] Prompt queued while session is busy: size=${promptQueue.size()}/${MAX_QUEUED_PROMPTS}`, ); - await replyWithKeyboard(ctx, t("queue.added")); + await replyWithKeyboard( + ctx, + t("queue.added", { count: String(promptQueue.size()), max: String(MAX_QUEUED_PROMPTS) }), + ); return true; } diff --git a/src/bot/handlers/prompt.ts b/src/bot/handlers/prompt.ts index f5da1cc6a..e0195aeeb 100644 --- a/src/bot/handlers/prompt.ts +++ b/src/bot/handlers/prompt.ts @@ -1,7 +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 { randomUUID } from "node:crypto"; +import { createOpencodeMessageId } from "../../utils/opencode-message-id.js"; import { opencodeClient } from "../../opencode/client.js"; import { clearSession, @@ -62,8 +62,14 @@ export { let botInstance: Bot | null = null; let chatIdInstance: number | null = null; +export type PromptTarget = { + sessionId: string | null; + directory: string; +}; + type ProcessPromptOptions = { responseMode?: PromptResponseMode; + target?: PromptTarget; }; export function getPromptBotInstance(): Bot | null { @@ -194,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.`, @@ -370,7 +387,7 @@ export async function processUserPrompt( configuredProviderID: storedModel.providerID, configuredModelID: storedModel.modelID, }); - const promptMessageId = randomUUID(); + const promptMessageId = createOpencodeMessageId(); setPromptResponseMode(currentSession.id, promptMessageId, responseMode); externalUserInputSuppressionManager.registerMessage(currentSession.id, promptMessageId); diff --git a/src/bot/menus/inline-menu.ts b/src/bot/menus/inline-menu.ts index 653387864..028747cf8 100644 --- a/src/bot/menus/inline-menu.ts +++ b/src/bot/menus/inline-menu.ts @@ -104,6 +104,10 @@ export async function replyWithInlineMenu( options: InlineMenuReplyOptions, ): 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", diff --git a/src/bot/messages/external-user-input-notification.ts b/src/bot/messages/external-user-input-notification.ts index 8f60dba10..78bd0b041 100644 --- a/src/bot/messages/external-user-input-notification.ts +++ b/src/bot/messages/external-user-input-notification.ts @@ -46,7 +46,7 @@ export async function deliverExternalUserInputNotification({ return false; } - if (consumeSuppressedInput(sessionId, messageId, text)) { + if (consumeSuppressedInput(sessionId, messageId)) { return false; } diff --git a/src/bot/middleware/interaction-guard.ts b/src/bot/middleware/interaction-guard.ts index f1781f42c..1ae96ff46 100644 --- a/src/bot/middleware/interaction-guard.ts +++ b/src/bot/middleware/interaction-guard.ts @@ -5,7 +5,6 @@ import { reconcileForegroundBusyState } from "../../app/services/run-control-ser import { canQueueMediaPrompt, rejectQueuedMediaBeforePreparation, - dispatchNextQueuedPrompt, shouldSuggestPromptQueue, tryEnqueuePrompt, } from "../handlers/prompt-queue-dispatch.js"; diff --git a/src/bot/services/event-subscription-service.ts b/src/bot/services/event-subscription-service.ts index f9b1ea2f8..9a9f4f917 100644 --- a/src/bot/services/event-subscription-service.ts +++ b/src/bot/services/event-subscription-service.ts @@ -729,11 +729,10 @@ class EventSubscriptionService implements BotEventSubscriptionService { sessionId, messageId, text: messageText, - consumeSuppressedInput: (incomingSessionId, incomingMessageId, incomingText) => + consumeSuppressedInput: (incomingSessionId, incomingMessageId) => externalUserInputSuppressionManager.consumeMessage( incomingSessionId, incomingMessageId, - incomingText, ), }); } catch (err) { @@ -1176,8 +1175,6 @@ class EventSubscriptionService implements BotEventSubscriptionService { await this.sessionCompletionTasks.get(sessionId)?.catch(() => undefined); const completedRun = assistantRunState.finishRun(sessionId, "session_idle"); - clearPromptResponseMode(sessionId); - externalUserInputSuppressionManager.clearSession(sessionId); if (!this.botInstance || !this.chatIdInstance) { foregroundSessionState.markIdle(sessionId); diff --git a/src/i18n/ar.ts b/src/i18n/ar.ts index 748c77e18..4228364a7 100644 --- a/src/i18n/ar.ts +++ b/src/i18n/ar.ts @@ -68,7 +68,7 @@ export const ar: I18nDictionary = { "bot.thinking": "💭 جارٍ التفكير...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ جارٍ العمل", - "progress.compact.finished_header": "✅ اكتمل الرد", + "progress.compact.finished_header": "✅ اكتمل العمل", "progress.compact.thinking": "💭 جارٍ التفكير...", "progress.compact.responding": "✍️ جارٍ كتابة الرد...", "progress.compact.waiting_question": "❓ في انتظار إجابتك...", @@ -391,8 +391,8 @@ export const ar: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 الإعداد الافتراضي", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 جارٍ الإرسال إلى قائمة انتظار OpenCode.", - "queue.full": "⚠️ قائمة الإرسال المعلّق ممتلئة ({max}). احذف رسالة أو حاول مرة أخرى بعد قليل.", + "queue.added": "📥 أُضيفت إلى قائمة الانتظار ({count}/{max}). ستُرسل بعد انتهاء المهمة الحالية.", + "queue.full": "⚠️ قائمة الانتظار ممتلئة ({max}). احذف رسالة أو انتظر انتهاء المهمة الحالية.", "queue.media_limit": "⚠️ الوسائط في قائمة الانتظار محدودة بـ {maxSizeMb} MiB. انتظر إرسال عنصر ثم أعد المحاولة.", "queue.removed": "🗑 تمت إزالة الرسالة من قائمة الانتظار.", "queue.not_found": "لم تعد هذه الرسالة في قائمة الانتظار.", diff --git a/src/i18n/de.ts b/src/i18n/de.ts index f53b65308..62a6d9845 100644 --- a/src/i18n/de.ts +++ b/src/i18n/de.ts @@ -67,7 +67,7 @@ export const de: I18nDictionary = { "bot.thinking": "💭 Denke...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Arbeite", - "progress.compact.finished_header": "✅ Antwort abgeschlossen", + "progress.compact.finished_header": "✅ Arbeit abgeschlossen", "progress.compact.thinking": "💭 Denke...", "progress.compact.responding": "✍️ Schreibe Antwort...", "progress.compact.waiting_question": "❓ Warte auf deine Antwort...", @@ -420,10 +420,11 @@ export const de: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Standard", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 Wird an OpenCodes Warteschlange übermittelt.", + "queue.added": + "📥 Zur Warteschlange hinzugefügt ({count}/{max}). Die Nachricht wird gesendet, sobald die aktuelle Aufgabe abgeschlossen ist.", "queue.media_limit": "⚠️ Medien in der Warteschlange sind auf {maxSizeMb} MiB begrenzt. Warte, bis ein Eintrag gesendet wurde.", "queue.full": - "⚠️ Die Warteschlange ausstehender Übermittlungen ist voll ({max}). Entferne eine Nachricht oder versuche es gleich erneut.", + "⚠️ Die Warteschlange ist voll ({max}). Entferne eine Nachricht oder warte, bis die aktuelle Aufgabe abgeschlossen ist.", "queue.removed": "🗑 Nachricht aus der Warteschlange entfernt.", "queue.not_found": "Diese Nachricht ist nicht mehr in der Warteschlange.", "queue.disabled_hint": "Die Nachrichtenwarteschlange lässt sich in /settings aktivieren.", diff --git a/src/i18n/en.ts b/src/i18n/en.ts index d15de4105..6d4bae7ef 100644 --- a/src/i18n/en.ts +++ b/src/i18n/en.ts @@ -64,7 +64,7 @@ export const en = { "bot.thinking": "💭 Thinking...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Working", - "progress.compact.finished_header": "✅ Response complete", + "progress.compact.finished_header": "✅ Finished Work", "progress.compact.thinking": "💭 Thinking...", "progress.compact.responding": "✍️ Writing answer...", "progress.compact.waiting_question": "❓ Waiting for your answer...", @@ -401,9 +401,8 @@ export const en = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Default", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 Forwarding to OpenCode's queue.", - "queue.full": - "⚠️ The pending submission queue is full ({max}). Remove a message or try again shortly.", + "queue.added": "📥 Added to queue ({count}/{max}). It will be sent when the current task finishes.", + "queue.full": "⚠️ Queue is full ({max}). Remove a message or wait for the current task to finish.", "queue.media_limit": "⚠️ Queued media is limited to {maxSizeMb} MiB. Wait for an item to send, then try again.", "queue.removed": "🗑 Message removed from the queue.", "queue.not_found": "This message is no longer in the queue.", diff --git a/src/i18n/es.ts b/src/i18n/es.ts index eb3d5a32a..4ba8c672b 100644 --- a/src/i18n/es.ts +++ b/src/i18n/es.ts @@ -67,7 +67,7 @@ export const es: I18nDictionary = { "bot.thinking": "💭 Pensando...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Trabajando", - "progress.compact.finished_header": "✅ Respuesta completada", + "progress.compact.finished_header": "✅ Trabajo terminado", "progress.compact.thinking": "💭 Pensando...", "progress.compact.responding": "✍️ Escribiendo respuesta...", "progress.compact.waiting_question": "❓ Esperando tu respuesta...", @@ -417,10 +417,11 @@ export const es: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Predeterminado", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 Enviando a la cola de OpenCode.", + "queue.added": + "📥 Añadido a la cola ({count}/{max}). Se enviará cuando termine la tarea actual.", "queue.media_limit": "⚠️ Los archivos multimedia en cola están limitados a {maxSizeMb} MiB. Espera a que se envíe un elemento.", "queue.full": - "⚠️ La cola de envíos pendientes está llena ({max}). Elimina un mensaje o inténtalo de nuevo en breve.", + "⚠️ La cola está llena ({max}). Elimina un mensaje o espera a que termine la tarea actual.", "queue.removed": "🗑 Mensaje eliminado de la cola.", "queue.not_found": "Este mensaje ya no está en la cola.", "queue.disabled_hint": "La cola de mensajes se activa en /settings.", diff --git a/src/i18n/fr.ts b/src/i18n/fr.ts index 277c0074c..2c5a9bc57 100644 --- a/src/i18n/fr.ts +++ b/src/i18n/fr.ts @@ -67,7 +67,7 @@ export const fr: I18nDictionary = { "bot.thinking": "💭 Réflexion en cours...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Travail en cours", - "progress.compact.finished_header": "✅ Réponse terminée", + "progress.compact.finished_header": "✅ Travail terminé", "progress.compact.thinking": "💭 Réflexion en cours...", "progress.compact.responding": "✍️ Rédaction de la réponse...", "progress.compact.waiting_question": "❓ En attente de votre réponse...", @@ -421,10 +421,11 @@ export const fr: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Par défaut", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 Envoi vers la file d'attente d'OpenCode.", + "queue.added": + "📥 Ajouté à la file d'attente ({count}/{max}). Le message sera envoyé à la fin de la tâche en cours.", "queue.media_limit": "⚠️ Les médias en file sont limités à {maxSizeMb} MiB. Attendez l'envoi d'un élément.", "queue.full": - "⚠️ La file des envois en attente est pleine ({max}). Supprimez un message ou réessayez bientôt.", + "⚠️ La file d'attente est pleine ({max}). Supprimez un message ou attendez la fin de la tâche en cours.", "queue.removed": "🗑 Message retiré de la file d'attente.", "queue.not_found": "Ce message n'est plus dans la file d'attente.", "queue.disabled_hint": "La file d'attente des messages s'active dans /settings.", diff --git a/src/i18n/it.ts b/src/i18n/it.ts index d2774767c..53775c9d9 100644 --- a/src/i18n/it.ts +++ b/src/i18n/it.ts @@ -68,7 +68,7 @@ export const it: I18nDictionary = { "bot.thinking": "💭 Sto pensando...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Al lavoro", - "progress.compact.finished_header": "✅ Risposta completata", + "progress.compact.finished_header": "✅ Lavoro completato", "progress.compact.thinking": "💭 Sto pensando...", "progress.compact.responding": "✍️ Scrivo la risposta...", "progress.compact.waiting_question": "❓ In attesa della tua risposta...", @@ -416,8 +416,8 @@ export const it: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Predefinito", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 Invio alla coda di OpenCode.", - "queue.full": "⚠️ La coda degli invii in sospeso è piena ({max}). Rimuovi un messaggio o riprova tra poco.", + "queue.added": "📥 Aggiunto alla coda ({count}/{max}). Verrà inviato quando l'attività corrente termina.", + "queue.full": "⚠️ La coda è piena ({max}). Rimuovi un messaggio o attendi che l'attività corrente termini.", "queue.media_limit": "⚠️ I media in coda sono limitati a {maxSizeMb} MiB. Attendi l'invio di un elemento e riprova.", "queue.removed": "🗑 Messaggio rimosso dalla coda.", "queue.not_found": "Questo messaggio non è più in coda.", diff --git a/src/i18n/ko.ts b/src/i18n/ko.ts index c09139b10..7cbb9a821 100644 --- a/src/i18n/ko.ts +++ b/src/i18n/ko.ts @@ -73,7 +73,7 @@ export const ko: I18nDictionary = { "bot.thinking": "💭 생각하는 중...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ 작업 중", - "progress.compact.finished_header": "✅ 응답 완료", + "progress.compact.finished_header": "✅ 작업 완료", "progress.compact.thinking": "💭 생각하는 중...", "progress.compact.responding": "✍️ 답변 작성 중...", "progress.compact.waiting_question": "❓ 답변을 기다리는 중...", @@ -410,8 +410,8 @@ export const ko: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 기본값", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 OpenCode 대기열로 전달 중입니다.", - "queue.full": "⚠️ 제출 대기열이 가득 찼습니다 ({max}). 메시지를 삭제하거나 잠시 후 다시 시도해 주세요.", + "queue.added": "📥 대기열에 추가되었습니다 ({count}/{max}). 현재 작업이 끝나면 전송됩니다.", + "queue.full": "⚠️ 대기열이 가득 찼습니다 ({max}). 메시지를 삭제하거나 현재 작업이 끝날 때까지 기다려 주세요.", "queue.media_limit": "⚠️ 대기열 미디어는 총 {maxSizeMb} MiB로 제한됩니다. 항목이 전송된 후 다시 시도하세요.", "queue.removed": "🗑 대기열에서 메시지를 삭제했습니다.", "queue.not_found": "이 메시지는 더 이상 대기열에 없습니다.", diff --git a/src/i18n/pt.ts b/src/i18n/pt.ts index 0a459213a..5477ce8a3 100644 --- a/src/i18n/pt.ts +++ b/src/i18n/pt.ts @@ -66,7 +66,7 @@ export const pt: I18nDictionary = { "bot.thinking": "💭 Pensando...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Trabalhando", - "progress.compact.finished_header": "✅ Resposta concluída", + "progress.compact.finished_header": "✅ Trabalho concluído", "progress.compact.thinking": "💭 Pensando...", "progress.compact.responding": "✍️ Escrevendo resposta...", "progress.compact.waiting_question": "❓ Aguardando sua resposta...", @@ -418,10 +418,11 @@ export const pt: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Padrão", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 Enviando para a fila do OpenCode.", + "queue.added": + "📥 Adicionado à fila ({count}/{max}). Será enviado quando a tarefa atual terminar.", "queue.media_limit": "⚠️ A mídia na fila está limitada a {maxSizeMb} MiB. Aguarde o envio de um item.", "queue.full": - "⚠️ A fila de envios pendentes está cheia ({max}). Remova uma mensagem ou tente novamente em breve.", + "⚠️ A fila está cheia ({max}). Remova uma mensagem ou aguarde o término da tarefa atual.", "queue.removed": "🗑 Mensagem removida da fila.", "queue.not_found": "Esta mensagem não está mais na fila.", "queue.disabled_hint": "A fila de mensagens pode ser ativada em /settings.", diff --git a/src/i18n/ru.ts b/src/i18n/ru.ts index 9e6321301..f7ce7fd8d 100644 --- a/src/i18n/ru.ts +++ b/src/i18n/ru.ts @@ -64,7 +64,7 @@ export const ru: I18nDictionary = { "bot.thinking": "💭 Думаю...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ Работаю", - "progress.compact.finished_header": "✅ Ответ завершён", + "progress.compact.finished_header": "✅ Работа завершена", "progress.compact.thinking": "💭 Думаю...", "progress.compact.responding": "✍️ Пишу ответ...", "progress.compact.waiting_question": "❓ Жду ваш ответ...", @@ -404,8 +404,8 @@ export const ru: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 Default", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 Отправляется в очередь OpenCode.", - "queue.full": "⚠️ Очередь ожидающих отправок заполнена ({max}). Удалите сообщение или повторите попытку чуть позже.", + "queue.added": "📥 Добавлено в очередь ({count}/{max}). Сообщение уйдёт после завершения текущей задачи.", + "queue.full": "⚠️ Очередь заполнена ({max}). Удалите сообщение или дождитесь завершения текущей задачи.", "queue.media_limit": "⚠️ Медиа в очереди ограничены {maxSizeMb} MiB. Дождитесь отправки элемента и повторите попытку.", "queue.removed": "🗑 Сообщение удалено из очереди.", "queue.not_found": "Этого сообщения больше нет в очереди.", diff --git a/src/i18n/zh.ts b/src/i18n/zh.ts index e4fa8420f..7678bba03 100644 --- a/src/i18n/zh.ts +++ b/src/i18n/zh.ts @@ -59,7 +59,7 @@ export const zh: I18nDictionary = { "bot.thinking": "💭 思考中...", "progress.compact.activity": "{header}\n{activity}", "progress.compact.working_header": "⏳ 工作中", - "progress.compact.finished_header": "✅ 回复完成", + "progress.compact.finished_header": "✅ 工作完成", "progress.compact.thinking": "💭 思考中...", "progress.compact.responding": "✍️ 正在撰写回复...", "progress.compact.waiting_question": "❓ 等待你的回答...", @@ -368,8 +368,8 @@ export const zh: I18nDictionary = { "keyboard.variant": "💭 {name}", "keyboard.variant_default": "💡 默认", "keyboard.queued_prompt": "❌ {index}. {text}", - "queue.added": "📥 正在提交到 OpenCode 队列。", - "queue.full": "⚠️ 待提交队列已满({max})。请删除一条消息或稍后重试。", + "queue.added": "📥 已加入队列({count}/{max})。当前任务完成后将自动发送。", + "queue.full": "⚠️ 队列已满({max})。请删除一条消息或等待当前任务完成。", "queue.media_limit": "⚠️ 队列媒体总大小限制为 {maxSizeMb} MiB。请等待一个项目发送后重试。", "queue.removed": "🗑 消息已从队列中移除。", "queue.not_found": "该消息已不在队列中。", 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 5864fd6a9..e7fcb41b7 100644 --- a/tests/app/managers/external-input-suppression-manager.test.ts +++ b/tests/app/managers/external-input-suppression-manager.test.ts @@ -10,33 +10,6 @@ describe("external-input/suppression", () => { vi.useRealTimers(); }); - it("consumes a matching suppressed input for the same session", () => { - externalUserInputSuppressionManager.register("session-1", "Review README"); - - expect(externalUserInputSuppressionManager.consume("session-1", "Review README")).toBe(true); - expect(externalUserInputSuppressionManager.consume("session-1", "Review README")).toBe(false); - }); - - it("does not consume a suppressed input from another session", () => { - externalUserInputSuppressionManager.register("session-1", "Review README"); - - expect(externalUserInputSuppressionManager.consume("session-2", "Review README")).toBe(false); - }); - - it("does not consume different text", () => { - externalUserInputSuppressionManager.register("session-1", "Review README"); - - expect(externalUserInputSuppressionManager.consume("session-1", "Review tests")).toBe(false); - }); - - it("expires stale suppression entries", () => { - externalUserInputSuppressionManager.register("session-1", "Review README", 1_000); - - expect(externalUserInputSuppressionManager.consume("session-1", "Review README", 61_001)).toBe( - false, - ); - }); - it("retires an unobserved admitted identity after the bounded lifetime", () => { vi.useFakeTimers(); externalUserInputSuppressionManager.registerMessage("session-1", "message-1"); @@ -47,11 +20,7 @@ describe("external-input/suppression", () => { expect(externalUserInputSuppressionManager.__getMessageCountForTests()).toBe(0); expect( - externalUserInputSuppressionManager.consumeMessage( - "session-1", - "message-1", - "different rendered text", - ), + externalUserInputSuppressionManager.consumeMessage("session-1", "message-1"), ).toBe(false); }); @@ -63,7 +32,7 @@ describe("external-input/suppression", () => { vi.advanceTimersByTime(2 * 60 * 1_000); expect( - externalUserInputSuppressionManager.consumeMessage("session-1", "message-1", "different"), + externalUserInputSuppressionManager.consumeMessage("session-1", "message-1"), ).toBe(true); expect(externalUserInputSuppressionManager.__getMessageCountForTests()).toBe(0); }); @@ -72,20 +41,7 @@ describe("external-input/suppression", () => { externalUserInputSuppressionManager.registerMessage("session-1", "message-1"); expect( - externalUserInputSuppressionManager.consumeMessage("session-1", "external-message", "same"), - ).toBe(false); - }); - - it("does not fall back to text suppression when an identity differs", () => { - externalUserInputSuppressionManager.register("session-1", "same text"); - externalUserInputSuppressionManager.registerMessage("session-1", "message-1"); - - expect( - externalUserInputSuppressionManager.consumeMessage( - "session-1", - "external-message", - "same text", - ), + externalUserInputSuppressionManager.consumeMessage("session-1", "external-message"), ).toBe(false); }); @@ -93,10 +49,10 @@ describe("external-input/suppression", () => { externalUserInputSuppressionManager.registerMessage("session-1", "message-1"); expect( - externalUserInputSuppressionManager.consumeMessage("session-2", "message-1", "unrelated"), + externalUserInputSuppressionManager.consumeMessage("session-2", "message-1"), ).toBe(false); expect( - externalUserInputSuppressionManager.consumeMessage("session-1", "message-1", "original"), + externalUserInputSuppressionManager.consumeMessage("session-1", "message-1"), ).toBe(true); }); diff --git a/tests/app/managers/telegram-input-order-manager.test.ts b/tests/app/managers/telegram-input-order-manager.test.ts index 436c5027a..56b0403c3 100644 --- a/tests/app/managers/telegram-input-order-manager.test.ts +++ b/tests/app/managers/telegram-input-order-manager.test.ts @@ -30,16 +30,4 @@ describe("app/managers/telegram-input-order-manager", () => { await expect(telegramInputOrderManager.waitForEarlier(777, 11)).resolves.toBeUndefined(); }); - it("holds a callback without a message id until the album is released", async () => { - telegramInputOrderManager.defer(777, 10); - const released = vi.fn(); - const waiting = telegramInputOrderManager.waitForPending(777).then(released); - - await Promise.resolve(); - expect(released).not.toHaveBeenCalled(); - - telegramInputOrderManager.release(777, 10); - await waiting; - expect(released).toHaveBeenCalledTimes(1); - }); }); diff --git a/tests/bot/handlers/prompt.test.ts b/tests/bot/handlers/prompt.test.ts index 70876b267..0eb896cee 100644 --- a/tests/bot/handlers/prompt.test.ts +++ b/tests/bot/handlers/prompt.test.ts @@ -278,7 +278,7 @@ describe("bot/handlers/prompt", () => { modelID: "gpt-5", }, variant: "default", - messageID: expect.any(String), + messageID: expect.stringMatching(/^msg_/), }); expect(mocked.sessionPromptMock).not.toHaveBeenCalled(); }); @@ -300,7 +300,7 @@ describe("bot/handlers/prompt", () => { ); }); - it("retains delivery state when promptAsync has an ambiguous transport failure", async () => { + it("notifies the user while retaining delivery state after an ambiguous transport failure", async () => { const ctx = createContext(); const deps = createDeps(); @@ -316,7 +316,10 @@ describe("bot/handlers/prompt", () => { backgroundTask.onError?.(error); }); - expect(deps.bot.api.sendMessage).not.toHaveBeenCalled(); + expect(deps.bot.api.sendMessage).toHaveBeenCalledWith( + 777, + "Failed to send request to OpenCode.", + ); expect(mocked.suppressionDiscardMock).not.toHaveBeenCalled(); }); diff --git a/tests/bot/messages/external-user-input-notification.test.ts b/tests/bot/messages/external-user-input-notification.test.ts index de5facabc..8f96ec01e 100644 --- a/tests/bot/messages/external-user-input-notification.test.ts +++ b/tests/bot/messages/external-user-input-notification.test.ts @@ -76,7 +76,6 @@ describe("bot/messages/external-user-input-notification", () => { expect(consumeSuppressedInput).toHaveBeenCalledWith( "session-1", "message-1", - "Review the parser", ); expect(mocked.sendBotTextMock).not.toHaveBeenCalled(); }); diff --git a/tests/bot/routers/telegram-input-order-routing.test.ts b/tests/bot/routers/telegram-input-order-routing.test.ts index 24364569a..1c7559d45 100644 --- a/tests/bot/routers/telegram-input-order-routing.test.ts +++ b/tests/bot/routers/telegram-input-order-routing.test.ts @@ -133,7 +133,7 @@ describe("bot/routers Telegram album ordering", () => { ); }); - it("holds a project-switch callback behind an earlier album", async () => { + 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); @@ -195,6 +195,6 @@ describe("bot/routers Telegram album ordering", () => { } as Update); await callback; - expect(order).toEqual(["album:session-a", "project:session-b"]); + expect(order).toEqual(["project:session-b"]); }); }); diff --git a/tests/bot/streaming/compact-progress-streamer.test.ts b/tests/bot/streaming/compact-progress-streamer.test.ts index 10bbe6581..6d9dfa920 100644 --- a/tests/bot/streaming/compact-progress-streamer.test.ts +++ b/tests/bot/streaming/compact-progress-streamer.test.ts @@ -21,7 +21,7 @@ describe("bot/streaming/compact-progress-streamer", () => { expect(editText).toHaveBeenCalledWith( "s1", 10, - "✅ Response complete\ntool calls: 0 · changed files: 0", + "✅ Finished Work\ntool calls: 0 · changed files: 0", ); }); @@ -67,7 +67,7 @@ describe("bot/streaming/compact-progress-streamer", () => { expect(editText).toHaveBeenCalledWith( "s1", 10, - "✅ Response complete\ntool calls: 0 · changed files: 0", + "✅ Finished Work\ntool calls: 0 · changed files: 0", ); }); @@ -104,7 +104,7 @@ describe("bot/streaming/compact-progress-streamer", () => { 2, "s1", 10, - "✅ Response complete\ntool calls: 1 · changed files: 0", + "✅ Finished Work\ntool calls: 1 · changed files: 0", ); }); @@ -126,7 +126,7 @@ describe("bot/streaming/compact-progress-streamer", () => { expect(sendText).toHaveBeenCalledTimes(1); expect(sendText).toHaveBeenCalledWith( "s1", - "✅ Response complete\ntool calls: 2 · changed files: 2", + "✅ Finished Work\ntool calls: 2 · changed files: 2", ); expect(editText).not.toHaveBeenCalled(); });