repositories / will
will
owned by admin
src/agent/telegram-runtime.ts
Rawimport 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<void> {
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<MediaRef | undefined> {
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<IngestResultSummary> {
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<void> {
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<string, string> = {
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'
}
}