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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion PRODUCT.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`)
Expand Down
120 changes: 76 additions & 44 deletions src/app/managers/external-input-suppression-manager.ts
Original file line number Diff line number Diff line change
@@ -1,74 +1,106 @@
const SUPPRESSION_TTL_MS = 60_000;

interface SuppressionEntry {
text: string;
createdAt: number;
}

function normalizeExternalUserInputText(text: string): string {
return text.replace(/\r\n/g, "\n").trim();
}
const MESSAGE_SUPPRESSION_TTL_MS = 5 * 60_000;

class ExternalUserInputSuppressionManager {
private entriesBySession = new Map<string, SuppressionEntry[]>();
private messageIdsBySession = new Map<string, Map<string, number>>();
private messagePruneTimer: ReturnType<typeof setTimeout> | null = null;

register(sessionId: string, text: string, now: number = Date.now()): void {
const normalizedText = normalizeExternalUserInputText(text);
if (!sessionId || !normalizedText) {
registerMessage(sessionId: string, messageId: string, now: number = Date.now()): void {
if (!sessionId || !messageId) {
return;
}

this.prune(now);

const sessionEntries = this.entriesBySession.get(sessionId) ?? [];
sessionEntries.push({ text: normalizedText, createdAt: now });
this.entriesBySession.set(sessionId, sessionEntries);
this.pruneMessages(now);
const messageIds = this.messageIdsBySession.get(sessionId) ?? new Map<string, number>();
messageIds.set(messageId, now + MESSAGE_SUPPRESSION_TTL_MS);
this.messageIdsBySession.set(sessionId, messageIds);
this.scheduleMessagePrune(now);
}

consume(sessionId: string, text: string, now: number = Date.now()): boolean {
const normalizedText = normalizeExternalUserInputText(text);
if (!sessionId || !normalizedText) {
return false;
consumeMessage(sessionId: string, messageId: string): boolean {
this.pruneMessages(Date.now());
const messageIds = this.messageIdsBySession.get(sessionId);
if (messageIds?.delete(messageId)) {
if (messageIds.size === 0) {
this.messageIdsBySession.delete(sessionId);
}
this.scheduleMessagePrune(Date.now());
return true;
}

this.prune(now);

const sessionEntries = this.entriesBySession.get(sessionId);
if (!sessionEntries?.length) {
return false;
}
return false;
}

const entryIndex = sessionEntries.findIndex((entry) => entry.text === normalizedText);
if (entryIndex < 0) {
return false;
discardMessage(sessionId: string, messageId: string): void {
const messageIds = this.messageIdsBySession.get(sessionId);
if (!messageIds) {
return;
}

sessionEntries.splice(entryIndex, 1);
if (sessionEntries.length === 0) {
this.entriesBySession.delete(sessionId);
messageIds.delete(messageId);
if (messageIds.size === 0) {
this.messageIdsBySession.delete(sessionId);
}
this.scheduleMessagePrune(Date.now());
}

return true;
clearSession(sessionId: string): void {
this.messageIdsBySession.delete(sessionId);
this.scheduleMessagePrune(Date.now());
}

clearAll(): void {
this.entriesBySession.clear();
this.messageIdsBySession.clear();
if (this.messagePruneTimer) {
clearTimeout(this.messagePruneTimer);
this.messagePruneTimer = null;
}
}

__resetForTests(): void {
this.clearAll();
}

private prune(now: number): void {
for (const [sessionId, sessionEntries] of this.entriesBySession.entries()) {
const activeEntries = sessionEntries.filter((entry) => now - entry.createdAt <= SUPPRESSION_TTL_MS);
if (activeEntries.length === 0) {
this.entriesBySession.delete(sessionId);
continue;
__getMessageCountForTests(): number {
return Array.from(this.messageIdsBySession.values()).reduce(
(count, messageIds) => count + messageIds.size,
0,
);
}

private pruneMessages(now: number): void {
for (const [sessionId, messageIds] of this.messageIdsBySession.entries()) {
for (const [messageId, expiresAt] of messageIds.entries()) {
if (expiresAt <= now) {
messageIds.delete(messageId);
}
}
if (messageIds.size === 0) {
this.messageIdsBySession.delete(sessionId);
}
}
}

this.entriesBySession.set(sessionId, activeEntries);
private scheduleMessagePrune(now: number): void {
if (this.messagePruneTimer) {
clearTimeout(this.messagePruneTimer);
this.messagePruneTimer = null;
}

const expirations = Array.from(this.messageIdsBySession.values()).flatMap((messageIds) =>
Array.from(messageIds.values()),
);
if (expirations.length === 0) {
return;
}

const expiresAt = Math.min(...expirations);
this.messagePruneTimer = setTimeout(() => {
this.messagePruneTimer = null;
const pruneAt = Date.now();
this.pruneMessages(pruneAt);
this.scheduleMessagePrune(pruneAt);
}, Math.max(0, expiresAt - now));
this.messagePruneTimer.unref?.();
}
}

Expand Down
26 changes: 26 additions & 0 deletions src/app/managers/interaction-manager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down
35 changes: 35 additions & 0 deletions src/app/managers/prompt-response-mode-manager.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
export type PromptResponseMode = "text_only" | "text_and_tts";

const promptResponseModes = new Map<string, Map<string, PromptResponseMode>>();

export function setPromptResponseMode(
sessionId: string,
messageId: string,
responseMode: PromptResponseMode,
): void {
const modes = promptResponseModes.get(sessionId) ?? new Map<string, PromptResponseMode>();
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;
}
2 changes: 2 additions & 0 deletions src/app/managers/summary-aggregation-manager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ export interface SummaryInfo {
}

export interface MessageCompletionInfo {
parentMessageId?: string | undefined;
agent?: string | undefined;
providerID?: string | undefined;
modelID?: string | undefined;
Expand Down Expand Up @@ -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,
Expand Down
103 changes: 103 additions & 0 deletions src/app/managers/telegram-input-order-manager.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
import type { Context, NextFunction } from "grammy";

interface PendingWaiter {
messageId: number;
resolve: () => void;
}

class TelegramInputOrderManager {
private readonly pendingByChat = new Map<number, Set<number>>();
private readonly waitersByChat = new Map<number, PendingWaiter[]>();

defer(chatId: number, messageId: number): void {
const pending = this.pendingByChat.get(chatId) ?? new Set<number>();
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<void> {
if (!this.hasEarlierPending(chatId, messageId)) {
return Promise.resolve();
}

return new Promise<void>((resolve) => {
const waiters = this.waitersByChat.get(chatId) ?? [];
waiters.push({ messageId, resolve });
this.waitersByChat.set(chatId, waiters);
});
}

__resetForTests(): void {
for (const waiters of this.waitersByChat.values()) {
for (const waiter of waiters) {
waiter.resolve();
}
}
this.pendingByChat.clear();
this.waitersByChat.clear();
}

private hasEarlierPending(chatId: number, messageId: number): boolean {
return Array.from(this.pendingByChat.get(chatId) ?? []).some(
(pendingMessageId) => pendingMessageId < messageId,
);
}

private resolveReadyWaiters(chatId: number): void {
const waiters = this.waitersByChat.get(chatId);
if (!waiters) {
return;
}

const blocked: PendingWaiter[] = [];
for (const waiter of waiters) {
if (this.hasEarlierPending(chatId, waiter.messageId)) {
blocked.push(waiter);
} else {
waiter.resolve();
}
}

if (blocked.length > 0) {
this.waitersByChat.set(chatId, blocked);
} else {
this.waitersByChat.delete(chatId);
}
}
}

export const telegramInputOrderManager = new TelegramInputOrderManager();

export async function telegramInputOrderMiddleware(
ctx: Context,
next: NextFunction,
): Promise<void> {
const message = ctx.message;
const chatId = ctx.chat?.id;
if (chatId === undefined) {
await next();
return;
}

if (!message) {
await next();
return;
}

if (message.media_group_id) {
await next();
return;
}

await telegramInputOrderManager.waitForEarlier(chatId, message.message_id);
await next();
}
5 changes: 5 additions & 0 deletions src/app/services/external-user-input-service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,11 @@ export interface ExternalUserInputNotification {
rawFallbackText: string;
}

export type ConsumeSuppressedInput = (
sessionId: string,
messageId: string,
) => boolean;

function normalizeExternalUserInputText(text: string): string {
return text.replace(/\r\n/g, "\n").trim();
}
Expand Down
Loading
Loading