import { z } from "zod"; import { canonicalJson, sha256, type JsonObject } from "../core/json.js"; import { newId } from "../core/ids.js"; import type { EventCandidate, ThoughtEvent } from "../events/types.js"; import type { JazzThoughtStore } from "../jazz/store.js"; import type { SourceCursor } from "../store/types.js"; import { projectTelegramFeedbackJudgments } from "../training/telegram-reactions.js"; import { downloadTelegramImage, imageAttachmentMetadata, isImageDocument, selectPhoto, TelegramImagePermanentError, } from "./telegram-images.js"; const TELEGRAM_WEBHOOK_REVISION = "telegram-bot-api-webhook-v1"; const COMPATIBLE_TELEGRAM_REVISIONS = new Set([ "telegram-bot-api-v1", "telegram-bot-api-v2", TELEGRAM_WEBHOOK_REVISION, ]); export const TELEGRAM_WEBHOOK_ALLOWED_UPDATES = ["message", "edited_message", "message_reaction"] as const; const TELEGRAM_CORRECTION_PREFIX = "/correct "; const TELEGRAM_CORRECTION_MAX_CHARS = 2_000; const userSchema = z.object({ id: z.number().int(), is_bot: z.boolean(), first_name: z.string(), last_name: z.string().optional(), username: z.string().optional(), }).passthrough(); const chatSchema = z.object({ id: z.number().int(), type: z.enum(["private", "group", "supergroup", "channel"]), title: z.string().optional(), username: z.string().optional(), first_name: z.string().optional(), last_name: z.string().optional(), }).passthrough(); const photoSizeSchema = z.object({ file_id: z.string(), file_unique_id: z.string(), width: z.number().int(), height: z.number().int(), file_size: z.number().int().nonnegative().optional(), }).passthrough(); const fileSchema = z.object({ file_id: z.string(), file_unique_id: z.string(), file_name: z.string().optional(), mime_type: z.string().optional(), file_size: z.number().int().nonnegative().optional(), }).passthrough(); const messageSchema = z.object({ message_id: z.number().int(), message_thread_id: z.number().int().optional(), from: userSchema.optional(), date: z.number().int().nonnegative(), edit_date: z.number().int().nonnegative().optional(), chat: chatSchema, text: z.string().optional(), caption: z.string().optional(), reply_to_message: z.object({ message_id: z.number().int() }).passthrough().optional(), photo: z.array(photoSizeSchema).optional(), document: fileSchema.optional(), audio: fileSchema.optional(), voice: fileSchema.optional(), video: fileSchema.optional(), }).passthrough(); const reactionTypeSchema = z.object({ type: z.string().min(1), emoji: z.string().optional(), custom_emoji_id: z.string().optional(), }).passthrough(); const messageReactionSchema = z.object({ chat: chatSchema, message_id: z.number().int(), user: userSchema.optional(), actor_chat: chatSchema.optional(), date: z.number().int().nonnegative(), old_reaction: z.array(reactionTypeSchema), new_reaction: z.array(reactionTypeSchema), }).passthrough(); const updateSchema = z.object({ update_id: z.number().int().nonnegative(), message: messageSchema.optional(), edited_message: messageSchema.optional(), message_reaction: messageReactionSchema.optional(), }).passthrough(); const envelopeSchema = z.object({ ok: z.boolean(), result: z.unknown().optional(), error_code: z.number().int().optional(), description: z.string().optional(), parameters: z.object({ retry_after: z.number().int().positive().optional() }).passthrough().optional(), }).passthrough(); const webhookInfoSchema = z.object({ url: z.string(), has_custom_certificate: z.boolean(), pending_update_count: z.number().int().nonnegative(), ip_address: z.string().optional(), last_error_date: z.number().int().nonnegative().optional(), last_error_message: z.string().optional(), last_synchronization_error_date: z.number().int().nonnegative().optional(), max_connections: z.number().int().positive().optional(), allowed_updates: z.array(z.string()).optional(), }).passthrough(); const botCommandSchema = z.object({ command: z.string().regex(/^[a-z0-9_]{1,32}$/), description: z.string().min(1).max(256), }); const botCommandsSchema = z.array(botCommandSchema).max(100); const menuButtonSchema = z.object({ type: z.string().min(1) }).passthrough(); const botNameSchema = z.object({ name: z.string().min(1).max(64) }).passthrough(); export type TelegramBotUser = z.infer; export type TelegramBotUpdate = z.infer; export type TelegramWebhookInfo = z.infer; export type TelegramBotCommand = z.infer; export interface TelegramBotClientOptions { token: string; baseUrl?: string; requestTimeoutMs?: number; } export interface TelegramSendResult { messageId: string; chatId: string; date: string; } export class TelegramBotApiError extends Error { constructor( message: string, readonly errorCode?: number, readonly retryAfterSeconds?: number, ) { super(message); this.name = "TelegramBotApiError"; } } export class TelegramBotClient { private readonly token: string; private readonly baseUrl: string; private readonly requestTimeoutMs: number; private identityValue?: TelegramBotUser; constructor(options: TelegramBotClientOptions) { this.token = required(options.token, "Telegram bot token"); this.baseUrl = (options.baseUrl ?? "https://api.telegram.org").replace(/\/$/, ""); this.requestTimeoutMs = boundedPositiveInteger(options.requestTimeoutMs ?? 35_000, "requestTimeoutMs", 120_000); } async identity(signal?: AbortSignal): Promise { this.identityValue ??= await this.call("getMe", {}, userSchema, signal); return this.identityValue; } async setWebhook(options: { url: string; secretToken: string; dropPendingUpdates?: boolean; signal?: AbortSignal; }): Promise { await this.call("setWebhook", { url: required(options.url, "Telegram webhook URL"), secret_token: webhookSecret(options.secretToken), max_connections: 1, allowed_updates: [...TELEGRAM_WEBHOOK_ALLOWED_UPDATES], drop_pending_updates: options.dropPendingUpdates ?? false, }, z.literal(true), options.signal); } async deleteWebhook(options: { dropPendingUpdates?: boolean; signal?: AbortSignal } = {}): Promise { await this.call("deleteWebhook", { drop_pending_updates: options.dropPendingUpdates ?? false, }, z.literal(true), options.signal); } async getWebhookInfo(signal?: AbortSignal): Promise { return this.call("getWebhookInfo", {}, webhookInfoSchema, signal); } async setCommands(commands: TelegramBotCommand[], signal?: AbortSignal): Promise { const validated = botCommandsSchema.parse(commands); await this.call("setMyCommands", { commands: validated }, z.literal(true), signal); } async getCommands(signal?: AbortSignal): Promise { return this.call("getMyCommands", {}, botCommandsSchema, signal); } async setName(name: string, signal?: AbortSignal): Promise { await this.call("setMyName", { name: required(name, "Telegram bot name") }, z.literal(true), signal); } async getName(signal?: AbortSignal): Promise { return (await this.call("getMyName", {}, botNameSchema, signal)).name; } async setCommandsMenuButton(chatId: string, signal?: AbortSignal): Promise { await this.call("setChatMenuButton", { chat_id: required(chatId, "Telegram chat id"), menu_button: { type: "commands" }, }, z.literal(true), signal); } async getMenuButton(chatId: string, signal?: AbortSignal): Promise { const result = await this.call("getChatMenuButton", { chat_id: required(chatId, "Telegram chat id"), }, menuButtonSchema, signal); return result.type; } async sendMessage(chatId: string, text: string, signal?: AbortSignal): Promise { const message = await this.call("sendMessage", { chat_id: required(chatId, "Telegram chat id"), text: required(text, "Telegram message text"), link_preview_options: { is_disabled: true }, }, messageSchema, signal); return { messageId: String(message.message_id), chatId: String(message.chat.id), date: new Date(message.date * 1_000).toISOString(), }; } async sendTyping(chatId: string, signal?: AbortSignal): Promise { const timeout = AbortSignal.timeout(3_000); const boundedSignal = signal ? AbortSignal.any([signal, timeout]) : timeout; await this.call("sendChatAction", { chat_id: required(chatId, "Telegram chat id"), action: "typing", }, z.literal(true), boundedSignal); } /** * Expose the base URL and token for image extraction in the webhook process. * Used by TelegramBotConnector to download admitted images via getFile. */ get imageExtractionConfig(): { baseUrl: string; token: string } { return { baseUrl: this.baseUrl, token: this.token }; } private async call( method: string, body: Record, resultSchema: z.ZodType, signal?: AbortSignal, ): Promise { const timeout = AbortSignal.timeout(this.requestTimeoutMs); const combinedSignal = signal ? AbortSignal.any([signal, timeout]) : timeout; let response: Response; try { response = await fetch(`${this.baseUrl}/bot${this.token}/${method}`, { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify(body), signal: combinedSignal, }); } catch (error) { const reason = error instanceof Error && error.name === "AbortError" ? "timed out or was cancelled" : "failed"; throw new TelegramBotApiError(`Telegram Bot API ${method} request ${reason}`); } let raw: unknown; try { raw = await response.json(); } catch { throw new TelegramBotApiError(`Telegram Bot API ${method} returned non-JSON HTTP ${response.status}`); } const envelope = envelopeSchema.parse(raw); if (!response.ok || !envelope.ok) { throw new TelegramBotApiError( envelope.description ?? `Telegram Bot API ${method} failed with HTTP ${response.status}`, envelope.error_code ?? response.status, envelope.parameters?.retry_after, ); } return resultSchema.parse(envelope.result); } } export interface TelegramReactionFeedbackRoute { chatId: string; allowedUserIds: string[]; } export interface TelegramBotConnectorOptions { id: string; bot: TelegramBotUser; allowedChatIds: string[]; reactionFeedback?: TelegramReactionFeedbackRoute[]; feedbackProjector?: typeof projectTelegramFeedbackJudgments; /** * Optional image extraction configuration. When provided, the connector * downloads at most one admitted PNG/JPEG image per message via getFile, * validates magic bytes, and stores it content-addressed under the artifact root. * The event carries only an opaque relative reference, never bytes or tokens. */ imageExtraction?: { client: TelegramBotClient; artifactRoot: string; }; } export interface TelegramWebhookIngestResult { status: "updated" | "unchanged"; accepted: number; ignored: number; inserted: number; unchanged: number; events: ThoughtEvent[]; cursor: SourceCursor; } export class TelegramBotConnector { readonly kind = "telegram" as const; readonly id: string; private readonly bot: TelegramBotUser; private readonly allowedChatIds: Set; private readonly reactionUsersByChat = new Map>(); private readonly imageExtraction?: { client: TelegramBotClient; artifactRoot: string } | undefined; private readonly feedbackProjector: typeof projectTelegramFeedbackJudgments; constructor(options: TelegramBotConnectorOptions) { this.id = required(options.id, "Telegram bot connector id"); this.bot = options.bot; this.allowedChatIds = new Set(options.allowedChatIds.map((value) => required(value, "Allowed Telegram chat id"))); if (this.allowedChatIds.size === 0) throw new Error("Telegram bot connector requires at least one allowed chat id"); for (const route of options.reactionFeedback ?? []) { const chatId = required(route.chatId, "Reaction feedback chat id"); if (!this.allowedChatIds.has(chatId)) throw new Error(`Reaction feedback chat is not an allowed Telegram chat: ${chatId}`); const userIds = new Set(route.allowedUserIds.map((value) => required(value, "Reaction feedback user id"))); if (userIds.size === 0) throw new Error(`Reaction feedback chat requires at least one allowed user: ${chatId}`); this.reactionUsersByChat.set(chatId, userIds); } this.imageExtraction = options.imageExtraction; this.feedbackProjector = options.feedbackProjector ?? projectTelegramFeedbackJudgments; } describe(): JsonObject { return { id: this.id, kind: this.kind, transport: "telegram-bot-api-webhook", revision: TELEGRAM_WEBHOOK_REVISION, botIdHash: sha256(String(this.bot.id)), allowedChatIdsHash: sha256([...this.allowedChatIds].sort().join("\n")), reactionFeedbackRoutesHash: sha256([...this.reactionUsersByChat.entries()] .sort(([left], [right]) => left.localeCompare(right)) .map(([chatId, userIds]) => `${chatId}:${[...userIds].sort().join(",")}`) .join("\n")), authority: "read-private-chat-updates-and-allowlisted-feedback", }; } async ingest(store: JazzThoughtStore, update: TelegramBotUpdate): Promise { const cursorId = `cursor:${this.id}`; const prior = await store.getSourceCursor(cursorId); const correlationId = `${this.id}:webhook-update:${update.update_id}`; const attemptId = newId("telegram_webhook"); const startedAt = new Date().toISOString(); await store.appendEvent(this.connectorEvent("started", attemptId, startedAt, { status: "started", transport: "telegram-bot-api-webhook", updateId: String(update.update_id), })); try { this.assertCursorCompatible(prior, String(this.bot.id)); const accepted = this.accepts(update); const feedbackUpdate = accepted && isFeedbackUpdate(update); const candidate = accepted ? await this.eventCandidate(store, update, this.bot, correlationId) : undefined; const candidates = candidate ? [candidate] : []; const completedAt = new Date().toISOString(); const cursor: SourceCursor = { id: cursorId, source: this.id, cursor: { revision: TELEGRAM_WEBHOOK_REVISION, botId: String(this.bot.id), highestUpdateId: Math.max(priorHighestUpdateId(prior), update.update_id), }, lastSuccessAt: completedAt, ...(prior?.lastFailureAt ? { lastFailureAt: prior.lastFailureAt } : {}), updatedAt: completedAt, }; const sourceCount = candidates.length; candidates.push(this.cursorEvent(attemptId, completedAt, cursor)); if (prior?.lastFailureAt && (!prior.lastSuccessAt || prior.lastFailureAt > prior.lastSuccessAt)) { candidates.push(this.connectorEvent("recovered", attemptId, completedAt, { status: "recovered", previousFailureAt: prior.lastFailureAt, })); } candidates.push(this.connectorEvent("completed", attemptId, completedAt, { status: accepted ? "updated" : "unchanged", updateId: String(update.update_id), accepted: accepted ? 1 : 0, ignored: accepted ? 0 : 1, })); const batch = await store.appendProducerBatch(candidates, cursor); const offeredEvents = batch.events.slice(0, sourceCount); const insertedIds = new Set(batch.inserted.map((event) => event.id)); const events = offeredEvents.filter((event) => insertedIds.has(event.id)); if (feedbackUpdate) await this.feedbackProjector(store, this.id); return { status: accepted ? "updated" : "unchanged", accepted: accepted ? 1 : 0, ignored: accepted ? 0 : 1, inserted: events.length, unchanged: sourceCount - events.length, events, cursor, }; } catch (error) { const failedAt = new Date().toISOString(); const message = error instanceof Error ? error.message : String(error); const durableCursor = await store.getSourceCursor(cursorId) ?? prior; const failedCursor: SourceCursor = { id: cursorId, source: this.id, cursor: durableCursor?.cursor ?? { revision: TELEGRAM_WEBHOOK_REVISION, botId: String(this.bot.id), highestUpdateId: -1, }, ...(durableCursor?.lastSuccessAt ? { lastSuccessAt: durableCursor.lastSuccessAt } : {}), lastFailureAt: failedAt, lastError: message, updatedAt: failedAt, }; await store.appendProducerBatch([this.connectorEvent("failed", attemptId, failedAt, { status: "failed", error: message, phase: "telegram-webhook-ingest", updateId: String(update.update_id), })], failedCursor); throw error; } } private accepts(update: TelegramBotUpdate): boolean { const message = update.edited_message ?? update.message; if (message) { if (!this.allowedChatIds.has(String(message.chat.id))) return false; if (!telegramCorrectionCommand(message.text)) return true; return message.chat.type === "private" && Boolean(message.from) && this.reactionUsersByChat.get(String(message.chat.id))?.has(String(message.from!.id)) === true; } const reaction = update.message_reaction; if (!reaction || reaction.chat.type !== "private" || !reaction.user) return false; return this.reactionUsersByChat.get(String(reaction.chat.id))?.has(String(reaction.user.id)) === true; } private async eventCandidate( store: JazzThoughtStore, update: TelegramBotUpdate, bot: TelegramBotUser, correlationId: string, ): Promise { const message = update.edited_message ?? update.message; if (!message) { return update.message_reaction ? this.reactionEventCandidate(store, update, bot, correlationId) : undefined; } const correction = telegramCorrectionCommand(message.text); if (correction) { return this.correctionEventCandidate(store, update, message, bot, correlationId, correction); } const edited = update.edited_message !== undefined; const occurredAt = new Date((message.edit_date ?? message.date) * 1_000).toISOString(); const senderId = message.from ? String(message.from.id) : String(message.chat.id); const senderName = message.from ? message.from.username ?? [message.from.first_name, message.from.last_name].filter(Boolean).join(" ") : message.chat.username ?? message.chat.title ?? "unknown"; const text = message.text ?? message.caption ?? ""; const attachments = await this.resolveAttachments(message); const payload = compact({ accountId: String(bot.id), accountUsername: bot.username, updateId: String(update.update_id), chatId: String(message.chat.id), messageId: String(message.message_id), senderId, senderName, chatType: message.chat.type === "private" ? "direct" : "channel", occurredAt: new Date(message.date * 1_000).toISOString(), editedAt: edited ? occurredAt : undefined, text, threadId: message.message_thread_id === undefined ? undefined : String(message.message_thread_id), replyToMessageId: message.reply_to_message?.message_id === undefined ? undefined : String(message.reply_to_message.message_id), attachments, transport: "telegram-bot-api-webhook", }); const conversationEligible = text.trim().length > 0 || attachments?.some((attachment) => ( attachment.kind === "image" && attachment.status === "stored" )) === true; return { type: conversationEligible ? "stream.thought.source.telegram.message" : "stream.thought.source.telegram.nonconversation", schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId: `${bot.id}:${message.chat.id}:${message.message_id}`, idempotencyKey: sha256(canonicalJson({ botId: bot.id, chatId: message.chat.id, messageId: message.message_id, revision: message.edit_date ?? "original", })), occurredAt, actor: senderId, correlationId, privacy: "sensitive", payload, }; } private async correctionEventCandidate( store: JazzThoughtStore, update: TelegramBotUpdate, message: z.infer, bot: TelegramBotUser, correlationId: string, command: TelegramCorrectionCommand, ): Promise { const chatId = String(message.chat.id); const messageId = String(message.message_id); const senderId = String(message.from!.id); const replyToMessageId = message.reply_to_message?.message_id === undefined ? undefined : String(message.reply_to_message.message_id); const resolution: CorrectionDeliveryResolution = command.status !== "valid" ? { status: command.status === "invalid" ? "invalid-command" : "replacement-too-long" } : !replyToMessageId ? { status: "missing-reply-target" } : await this.resolveCorrectionDelivery(store, chatId, replyToMessageId); const occurredAt = new Date((message.edit_date ?? message.date) * 1_000).toISOString(); const payload = compact({ accountId: String(bot.id), accountUsername: bot.username, updateId: String(update.update_id), chatId, messageId, senderId, chatType: "direct", occurredAt: new Date(message.date * 1_000).toISOString(), commandVersion: "telegram-correct-v1", replacementText: command.status === "valid" ? command.replacementText : undefined, replacementChars: command.replacementChars, replacementSha256: command.replacementSha256, replyToMessageId, resolutionStatus: resolution.status, deliveryReceiptEventId: resolution.deliveryReceiptEventId, runId: resolution.runId, outputEventId: resolution.outputEventId, sourceRootEventId: resolution.sourceRootEventId, transport: "telegram-bot-api-webhook", }); return { type: "stream.thought.source.telegram.correction", schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId: `${bot.id}:${message.chat.id}:${message.message_id}:correction`, idempotencyKey: sha256(canonicalJson({ botId: bot.id, chatId: message.chat.id, messageId: message.message_id, revision: message.edit_date ?? "original", kind: "correction", })), occurredAt, actor: senderId, ...(resolution.sourceRootEventId ? { rootEventId: resolution.sourceRootEventId } : {}), ...(resolution.deliveryReceiptEventId ? { parentEventId: resolution.deliveryReceiptEventId } : {}), correlationId: resolution.deliveryReceiptExternalId ?? correlationId, privacy: "sensitive", payload, }; } /** * Resolve message attachments. When image extraction is configured, admitted * PNG/JPEG images (photos or image documents) are downloaded via getFile, * validated, content-addressed, and stored under the artifact root. * The attachment metadata carries only an opaque relative path reference, * never raw bytes, base64, the bot token, file URLs, or absolute paths. */ private async resolveAttachments(message: z.infer): Promise { if (!this.imageExtraction) return attachmentsFor(message); const { client, artifactRoot } = this.imageExtraction; const { baseUrl, token } = client.imageExtractionConfig; const attachments: JsonObject[] = []; // Attempt to extract at most one admitted image from photo array if (message.photo && message.photo.length > 0) { const photo = selectPhoto(message.photo); if (photo) { try { const artifact = await downloadTelegramImage(photo.file_id, { baseUrl, token, artifactRoot, }); attachments.push({ ...imageAttachmentMetadata(artifact), id: photo.file_unique_id, reference: `telegram-file:${photo.file_id}`, width: photo.width, height: photo.height, }); } catch (error) { if (!(error instanceof TelegramImagePermanentError)) throw error; attachments.push(compact({ id: photo.file_unique_id, kind: "image", sizeBytes: photo.file_size, reference: `telegram-file:${photo.file_id}`, status: "rejected", reason: error.code, })); } } else { const smallest = [...message.photo].sort((left, right) => (left.width * left.height) - (right.width * right.height))[0]!; attachments.push(compact({ id: smallest.file_unique_id, kind: "image", sizeBytes: smallest.file_size, reference: `telegram-file:${smallest.file_id}`, status: "rejected", reason: "size-limit", })); } } // A Telegram photo array and image document are mutually exclusive in valid // updates. If a malformed update supplies both, the photo path owns the one // permitted image slot and the document is ignored. if ((!message.photo || message.photo.length === 0) && message.document && isImageDocument(message.document)) { try { const artifact = await downloadTelegramImage(message.document.file_id, { baseUrl, token, artifactRoot, }); attachments.push(compact({ ...imageAttachmentMetadata(artifact), id: message.document.file_unique_id, name: message.document.file_name, reference: `telegram-file:${message.document.file_id}`, })); } catch (error) { if (!(error instanceof TelegramImagePermanentError)) throw error; attachments.push(compact({ id: message.document.file_unique_id, name: message.document.file_name, kind: "image", sizeBytes: message.document.file_size, reference: `telegram-file:${message.document.file_id}`, status: "rejected", reason: error.code, })); } } // Add non-image attachments without extraction. for (const [kind, file] of [ ["file", message.document && !isImageDocument(message.document) ? message.document : undefined], ["audio", message.audio], ["audio", message.voice], ["video", message.video], ] as const) { if (!file) continue; attachments.push(compact({ id: file.file_unique_id, name: file.file_name, mimeType: file.mime_type, sizeBytes: file.file_size, kind, reference: `telegram-file:${file.file_id}`, })); } return attachments.length > 0 ? attachments : undefined; } private async reactionEventCandidate( store: JazzThoughtStore, update: TelegramBotUpdate, bot: TelegramBotUser, correlationId: string, ): Promise { const reaction = update.message_reaction!; const chatId = String(reaction.chat.id); const messageId = String(reaction.message_id); const senderId = String(reaction.user!.id); const resolution = await this.resolveReactionDelivery(store, chatId, messageId); const oldLabel = reactionLabel(reaction.old_reaction); const newLabel = reactionLabel(reaction.new_reaction); const feedbackAction = resolution.status === "resolved" ? newLabel && newLabel !== oldLabel ? "set" : oldLabel && !newLabel ? "retract" : "none" : "none"; const payload = compact({ accountId: String(bot.id), accountUsername: bot.username, updateId: String(update.update_id), chatId, messageId, senderId, chatType: "direct", occurredAt: new Date(reaction.date * 1_000).toISOString(), oldReactions: reaction.old_reaction.map(reactionDescriptor), newReactions: reaction.new_reaction.map(reactionDescriptor), mappingVersion: "telegram-thumbs-v1", resolutionStatus: resolution.status, feedbackAction, feedbackLabel: resolution.status === "resolved" ? newLabel : undefined, deliveryReceiptEventId: resolution.deliveryReceiptEventId, runId: resolution.runId, outputEventId: resolution.outputEventId, sourceRootEventId: resolution.sourceRootEventId, transport: "telegram-bot-api-webhook", }); return { type: "stream.thought.source.telegram.reaction", schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId: `${bot.id}:reaction:${update.update_id}`, idempotencyKey: sha256(canonicalJson({ botId: bot.id, updateId: update.update_id, kind: "message_reaction" })), occurredAt: new Date(reaction.date * 1_000).toISOString(), actor: senderId, ...(resolution.sourceRootEventId ? { rootEventId: resolution.sourceRootEventId } : {}), ...(resolution.deliveryReceiptEventId ? { parentEventId: resolution.deliveryReceiptEventId } : {}), correlationId: resolution.deliveryReceiptExternalId ?? correlationId, privacy: "sensitive", payload, }; } private async resolveCorrectionDelivery( store: JazzThoughtStore, chatId: string, messageId: string, ): Promise { const resolution = await this.resolveReactionDelivery(store, chatId, messageId); if (resolution.status !== "resolved") return resolution; const run = await store.getRun(resolution.runId!); if (!run || run.status !== "completed" || !run.result) return { status: "incomplete-run" }; if (run.outputEventIds.length === 0) return { status: "missing-output" }; if (run.outputEventIds.length !== 1) return { status: "ambiguous-output" }; const [output, trigger, delivery] = await Promise.all([ store.getEvent(run.outputEventIds[0]!), store.getEvent(run.triggerEventId), store.getEvent(resolution.deliveryReceiptEventId!), ]); if (!output) return { status: "missing-output" }; if ( !trigger || trigger.source !== this.id || trigger.payload.chatId !== chatId || output.type !== "stream.thought.derived.message.observation" || output.source !== `agent:${run.agentId}` || output.sourceKind !== "agent" || output.actor !== run.agentId || output.parentEventId !== trigger.id || output.rootEventId !== trigger.rootEventId || output.payload.runId !== run.id || output.payload.inputEventId !== trigger.id || !delivery || delivery.type !== "stream.thought.action.telegram.send.delivered" || delivery.source !== `telegram-dispatcher:${this.id}:${chatId}` || delivery.sourceKind !== "system" || delivery.actor !== delivery.source || delivery.parentEventId !== output.id || delivery.rootEventId !== trigger.rootEventId || delivery.payload.chatId !== chatId || !Array.isArray(delivery.payload.runIds) || delivery.payload.runIds.length !== 1 || delivery.payload.runIds[0] !== run.id || resolution.sourceRootEventId !== trigger.rootEventId ) return { status: "invalid-lineage" }; return { ...resolution, outputEventId: output.id, }; } private async resolveReactionDelivery( store: JazzThoughtStore, chatId: string, messageId: string, ): Promise { const matches = (await store.listEvents({ types: ["stream.thought.action.telegram.send.delivered"] })) .filter((event) => event.payload.chatId === chatId && event.payload.messageId === messageId); if (matches.length === 0) return { status: "unknown-delivery" }; if (matches.length !== 1) return { status: "ambiguous-delivery" }; const receipt = matches[0]!; const runIds = Array.isArray(receipt.payload.runIds) ? receipt.payload.runIds.filter((value): value is string => typeof value === "string") : []; if (runIds.length !== 1) return { status: "ambiguous-run" }; const run = await store.getRun(runIds[0]!); if (!run) return { status: "missing-run" }; const trigger = await store.getEvent(run.triggerEventId); if (!trigger || trigger.source !== this.id || trigger.payload.chatId !== chatId || receipt.source !== `telegram-dispatcher:${this.id}:${chatId}` || receipt.sourceKind !== "system" || receipt.actor !== receipt.source || receipt.rootEventId !== trigger.rootEventId) return { status: "invalid-lineage" }; return { status: "resolved", deliveryReceiptEventId: receipt.id, deliveryReceiptExternalId: receipt.externalId, runId: run.id, ...(run.outputEventIds[0] ? { outputEventId: run.outputEventIds[0] } : {}), sourceRootEventId: trigger.rootEventId, }; } private assertCursorCompatible(cursor: SourceCursor | undefined, botId: string): void { const revision = cursor?.cursor.revision; if (revision !== undefined && (typeof revision !== "string" || !COMPATIBLE_TELEGRAM_REVISIONS.has(revision))) { throw new Error("Telegram bot cursor revision does not match connector configuration"); } const priorBotId = cursor?.cursor.botId; if (priorBotId !== undefined && priorBotId !== botId) throw new Error("Telegram bot cursor belongs to a different bot identity"); priorHighestUpdateId(cursor); } private cursorEvent(correlationId: string, at: string, cursor: SourceCursor): EventCandidate { return { type: "stream.thought.connector.cursor.advanced", schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId: cursor.id, idempotencyKey: `${correlationId}:cursor`, occurredAt: at, actor: this.id, correlationId, privacy: "sensitive", payload: { source: this.id, cursorId: cursor.id, cursor: cursor.cursor }, }; } private connectorEvent( phase: "started" | "completed" | "failed" | "recovered", correlationId: string, at: string, payload: JsonObject, ): EventCandidate { const type = phase === "failed" ? "stream.thought.connector.failed" : phase === "recovered" ? "stream.thought.connector.recovered" : `stream.thought.connector.ingest.${phase}`; return { type, schemaVersion: 1, source: this.id, sourceKind: this.kind, externalId: correlationId, idempotencyKey: `${correlationId}:${phase}`, occurredAt: at, actor: this.id, correlationId, privacy: "sensitive", payload, }; } } function isFeedbackUpdate(update: TelegramBotUpdate): boolean { if (update.message_reaction) return true; const message = update.edited_message ?? update.message; return Boolean(message && telegramCorrectionCommand(message.text)); } interface ReactionDeliveryResolution { status: "resolved" | "unknown-delivery" | "ambiguous-delivery" | "ambiguous-run" | "missing-run" | "invalid-lineage"; deliveryReceiptEventId?: string; deliveryReceiptExternalId?: string; runId?: string; outputEventId?: string; sourceRootEventId?: string; } interface CorrectionDeliveryResolution { status: ReactionDeliveryResolution["status"] | "incomplete-run" | "ambiguous-output" | "missing-output" | "invalid-command" | "replacement-too-long" | "missing-reply-target"; deliveryReceiptEventId?: string; deliveryReceiptExternalId?: string; runId?: string; outputEventId?: string; sourceRootEventId?: string; } interface TelegramCorrectionCommand { status: "valid" | "invalid" | "too-long"; replacementText?: string; replacementChars: number; replacementSha256: string; } function telegramCorrectionCommand(text: string | undefined): TelegramCorrectionCommand | undefined { if (text === "/correct") { return { status: "invalid", replacementChars: 0, replacementSha256: sha256("") }; } if (!text?.startsWith(TELEGRAM_CORRECTION_PREFIX)) return undefined; const replacementText = text.slice(TELEGRAM_CORRECTION_PREFIX.length); if (replacementText.length === 0) { return { status: "invalid", replacementChars: 0, replacementSha256: sha256("") }; } if (replacementText.length > TELEGRAM_CORRECTION_MAX_CHARS) { return { status: "too-long", replacementChars: replacementText.length, replacementSha256: sha256(replacementText), }; } return { status: "valid", replacementText, replacementChars: replacementText.length, replacementSha256: sha256(replacementText), }; } function reactionLabel(reactions: Array>): "positive" | "negative" | undefined { if (reactions.length !== 1 || reactions[0]?.type !== "emoji") return undefined; if (reactions[0].emoji === "👍") return "positive"; if (reactions[0].emoji === "👎") return "negative"; return undefined; } function reactionDescriptor(reaction: z.infer): JsonObject { return compact({ type: reaction.type, emoji: reaction.emoji, customEmojiId: reaction.custom_emoji_id, }); } function attachmentsFor(message: z.infer): JsonObject[] | undefined { const attachments: JsonObject[] = []; const photo = message.photo?.at(-1); if (photo) attachments.push(compact({ id: photo.file_unique_id, kind: "image", sizeBytes: photo.file_size, reference: `telegram-file:${photo.file_id}`, })); for (const [kind, file] of [ ["file", message.document], ["audio", message.audio], ["audio", message.voice], ["video", message.video], ] as const) { if (!file) continue; attachments.push(compact({ id: file.file_unique_id, name: file.file_name, mimeType: file.mime_type, sizeBytes: file.file_size, kind, reference: `telegram-file:${file.file_id}`, })); } return attachments.length > 0 ? attachments : undefined; } function priorHighestUpdateId(cursor: SourceCursor | undefined): number { const revision = cursor?.cursor.revision; const value = revision === "telegram-bot-api-v1" || revision === "telegram-bot-api-v2" ? cursor?.cursor.updateOffset : cursor?.cursor.highestUpdateId; if (value === undefined) return -1; if (typeof value !== "number" || !Number.isSafeInteger(value) || value < (revision === TELEGRAM_WEBHOOK_REVISION ? -1 : 0)) { throw new Error("Telegram bot cursor update high-water mark must be a safe integer"); } return revision === "telegram-bot-api-v1" || revision === "telegram-bot-api-v2" ? Math.max(-1, value - 1) : value; } export function parseTelegramBotUpdate(value: unknown): TelegramBotUpdate { return updateSchema.parse(value); } function compact(value: Record): JsonObject { return Object.fromEntries(Object.entries(value).filter(([, child]) => child !== undefined)) as JsonObject; } function boundedPositiveInteger(value: number, label: string, maximum: number): number { if (!Number.isSafeInteger(value) || value <= 0 || value > maximum) throw new Error(`${label} must be a positive integer <= ${maximum}`); return value; } function required(value: string, label: string): string { const normalized = value.trim(); if (!normalized) throw new Error(`${label} is required`); return normalized; } function webhookSecret(value: string): string { const secret = required(value, "Telegram webhook secret"); if (!/^[A-Za-z0-9_-]{1,256}$/.test(secret)) { throw new Error("Telegram webhook secret must use 1-256 ASCII letters, digits, underscores, or hyphens"); } return secret; }