import type { AuditEvent, ChatMessageIn, MediaRef } from '../shared/events.ts' import { createLogger } from '../shared/log.ts' import { classifyInbound, mayHide, toTgMessage } from '../shared/telegram.ts' import { nowIso } from '../shared/time.ts' import type { WakeQueue } from './queue-store.ts' import type { ServiceClient } from './service-client.ts' import type { AuditSpool } from './spool.ts' import type { AgentState } from './state.ts' import type { TelegramClient } from './telegram-client.ts' const log = createLogger({ component: 'telegram' }) export interface TelegramRuntimeDeps { tg: TelegramClient service: ServiceClient queue: WakeQueue spool: AuditSpool state: AgentState chatId: number oliverId: string } export interface ChatWakePayload { messageId: number senderId?: string | undefined senderName?: string | undefined text?: string | undefined triggerReason: 'mention' | 'reply' | 'command' media: { digest: string; type: string; fileName?: string | undefined }[] replyToMessageId?: number | undefined } export interface IngestResultSummary { archived: number wakes: number hidden: number ignored: number } /** * Telegram ingestion loop (spec: archive every authorized delivered update; * only mention/reply/commands wake; /hide is deterministic, no model call). */ export function createTelegramRuntime(deps: TelegramRuntimeDeps) { const { tg, service, queue, spool, state, chatId, oliverId } = deps let stopped = false let polling = false const audit = (kind: string, payload: unknown) => void spool.append( { id: crypto.randomUUID(), ts: nowIso(), kind, payload } satisfies AuditEvent, {}, ) async function identifyBot(): Promise { const me = await tg.getMe() await state.update((d) => { d.botUserId = me.id d.botUsername = me.username }) log.info('bot identified', { username: me.username }) } async function downloadMedia( fileId: string, fallbackType: string, fileName?: string, ): Promise { const file = await tg.getFile(fileId) const bytes = await tg.downloadFile(file.filePath) const type = guessType(fileName, file.filePath, fallbackType) const meta = await service.putMedia(bytes, { type, ...(fileName !== undefined ? { originalName: fileName } : {}), telegramFileId: fileId, }) return { digest: meta.digest, type, size: meta.size, ...(fileName !== undefined ? { fileName } : {}), ...(fileId !== undefined ? { telegramFileId: fileId } : {}), } } /** Poll once; returns what happened so tests can drive single iterations. */ async function pollOnce(): Promise { const summary: IngestResultSummary = { archived: 0, wakes: 0, hidden: 0, ignored: 0 } const offset = state.data.telegramOffset ?? 0 const updates = await tg.getUpdates(offset, 0) for (const raw of updates) { const nextOffset = raw.update_id + 1 await state.update((d) => { d.telegramOffset = nextOffset }) const msg = toTgMessage(raw) if (!msg) continue if (msg.chatId !== chatId) { summary.ignored++ audit('telegram_ignored', { chatId: msg.chatId, updateId: msg.updateId }) continue } const botUserId = state.data.botUserId ?? -1 const botUsername = state.data.botUsername ?? '' const classification = classifyInbound(msg, botUserId, botUsername) // Media download before ingest so the archived message references digests. const mediaRefs: MediaRef[] = [] for (const m of msg.media ?? []) { try { const ref = await downloadMedia(m.fileId, kindToType(m.kind), m.fileName) if (ref) mediaRefs.push(ref) } catch (err) { log.warn('media download failed', { updateId: msg.updateId, error: err instanceof Error ? err.message : String(err), }) audit('telegram_media_failed', { updateId: msg.updateId, fileId: m.fileId }) } } const effectiveText = msg.text ?? msg.caption const chatRow: ChatMessageIn = { id: `tg-${msg.updateId}`, chatId: msg.chatId, messageId: msg.messageId, direction: 'in', ts: new Date(msg.date * 1000).toISOString(), ...(msg.updateId !== undefined ? { updateId: msg.updateId } : {}), ...(msg.senderId !== undefined ? { senderId: String(msg.senderId) } : {}), ...(msg.senderName !== undefined ? { senderName: msg.senderName } : {}), ...(msg.replyToMessageId !== undefined ? { replyToMessageId: msg.replyToMessageId } : {}), triggerReason: classification.reason, ...(effectiveText !== undefined ? { text: effectiveText } : {}), ...(mediaRefs.length > 0 ? { media: mediaRefs } : {}), } await service.ingestChat([chatRow]) summary.archived++ if (classification.hide) { const target = classification.hide.targetMessageId const targetRow = await service.chatByMessageId(target) const allowed = mayHide({ ...(targetRow?.senderId !== undefined ? { targetSenderId: targetRow.senderId } : {}), requesterId: chatRow.senderId ?? '', oliverId, }) if (allowed) { await service.hideMessage({ messageId: target, requesterId: chatRow.senderId ?? 'unknown', authority: chatRow.senderId === oliverId ? 'oliver' : 'sender', ts: nowIso(), }) summary.hidden++ audit('chat_hidden', { messageId: target, by: chatRow.senderId }) } else { audit('chat_hide_denied', { messageId: target, by: chatRow.senderId }) } continue } if (classification.wake) { const payload: ChatWakePayload = { messageId: msg.messageId, ...(msg.senderId !== undefined ? { senderId: String(msg.senderId) } : {}), ...(msg.senderName !== undefined ? { senderName: msg.senderName } : {}), ...(effectiveText !== undefined ? { text: effectiveText } : {}), triggerReason: classification.reason === 'none' ? 'command' : classification.reason, media: mediaRefs.map((r) => ({ digest: r.digest, type: r.type, ...(r.fileName !== undefined ? { fileName: r.fileName } : {}), })), ...(msg.replyToMessageId !== undefined ? { replyToMessageId: msg.replyToMessageId } : {}), } await queue.enqueue('chat', payload) summary.wakes++ } } return summary } async function run(): Promise { while (!stopped) { if (polling) return polling = true try { if (state.data.botUsername === undefined) await identifyBot() const summary = await pollOnce() if (summary.archived === 0 && summary.ignored === 0) { await new Promise((res) => setTimeout(res, 1000)) } } catch (err) { log.warn('poll failed', { error: err instanceof Error ? err.message : String(err) }) await new Promise((res) => setTimeout(res, 3000)) } finally { polling = false } } } return { identifyBot, pollOnce, run, stop() { stopped = true }, } } function guessType(fileName: string | undefined, filePath: string, fallback: string): string { const name = fileName ?? filePath const ext = name.includes('.') ? (name.split('.').pop() ?? '') : '' const map: Record = { png: 'image/png', jpg: 'image/jpeg', jpeg: 'image/jpeg', webp: 'image/webp', gif: 'image/gif', mp4: 'video/mp4', webm: 'video/webm', mp3: 'audio/mpeg', ogg: 'audio/ogg', oga: 'audio/ogg', wav: 'audio/wav', pdf: 'application/pdf', txt: 'text/plain', md: 'text/markdown', json: 'application/json', } return map[ext] ?? fallback } function kindToType(kind: string): string { switch (kind) { case 'photo': return 'image/jpeg' case 'animation': return 'image/gif' case 'video': case 'video_note': return 'video/mp4' case 'voice': return 'audio/ogg' case 'audio': return 'audio/mpeg' default: return 'application/octet-stream' } }