/** * Multi-session Telegram topic registry. * Stores runtime session bindings separately from durable Telegram bot config. */ import { createHash } from "node:crypto"; import { existsSync } from "node:fs"; import { mkdir, readFile, rename, writeFile } from "node:fs/promises"; import { homedir } from "node:os"; import { dirname, join, resolve } from "node:path"; export interface TelegramTopicBinding { chatId: number; threadId?: number; chatTitle?: string; topicName?: string; updatedAt: number; } export interface TelegramRegisteredSession { key: string; pid: number; cwd: string; sessionFile?: string; inboxPath: string; inboxCursor?: number; binding?: TelegramTopicBinding; bindCode?: string; bindCodeExpiresAt?: number; broadcastEnabled?: boolean; updatedAt: number; } export interface TelegramSessionRegistryState { sessions: Record; } export interface TelegramSessionIdentity { pid: number; cwd: string; sessionFile?: string; } export interface TelegramSessionRegistryOptions { agentDir?: string; registryPath?: string; now?: () => number; isPidAlive?: (pid: number) => boolean; } export interface TelegramSessionRegistry { getPath: () => string; getAgentDir: () => string; load: () => Promise; save: (state: TelegramSessionRegistryState) => Promise; register: (session: TelegramRegisteredSession) => Promise; unregister: (key: string) => Promise; getSession: (key: string) => Promise; findByTopic: (chatId: number, threadId: number | undefined) => Promise; setBindCode: (key: string, code: string, ttlMs: number) => Promise; bindTopicByCode: (code: string, binding: TelegramTopicBinding) => Promise; bindTopic: (key: string, binding: TelegramTopicBinding) => Promise; unbindTopic: (key: string) => Promise; setBroadcast: (key: string, enabled: boolean) => Promise; updateInboxCursor: (key: string, cursor: number) => Promise; } function defaultAgentDir(): string { return process.env.PI_CODING_AGENT_DIR ? resolve(process.env.PI_CODING_AGENT_DIR) : join(homedir(), ".pi", "agent"); } function defaultIsPidAlive(pid: number): boolean { try { process.kill(pid, 0); return true; } catch { return false; } } export function makeTelegramSessionKey(identity: TelegramSessionIdentity): string { return `${identity.pid}:${identity.cwd}:${identity.sessionFile ?? ""}`; } export function makeTelegramSessionInboxPath(agentDir: string, key: string): string { const digest = createHash("sha256").update(key).digest("hex"); return join(agentDir, "pi-telegram-router", "inbox", `${digest}.jsonl`); } export function normalizeBindCode(code: string): string { return code.trim().toUpperCase(); } export function pruneDeadTelegramSessions( state: TelegramSessionRegistryState, isPidAlive: (pid: number) => boolean, ): TelegramSessionRegistryState { const sessions: Record = {}; for (const [key, session] of Object.entries(state.sessions ?? {})) { if (isPidAlive(session.pid)) sessions[key] = session; } return { sessions }; } export function createTelegramSessionRegistry( options: TelegramSessionRegistryOptions = {}, ): TelegramSessionRegistry { const agentDir = options.agentDir ?? defaultAgentDir(); const registryPath = options.registryPath ?? join(agentDir, "pi-telegram-sessions.json"); const now = options.now ?? Date.now; const isPidAlive = options.isPidAlive ?? defaultIsPidAlive; async function load(): Promise { if (!existsSync(registryPath)) return { sessions: {} }; try { const parsed = JSON.parse(await readFile(registryPath, "utf8")) as TelegramSessionRegistryState; return pruneDeadTelegramSessions({ sessions: parsed.sessions ?? {} }, isPidAlive); } catch { const backupPath = `${registryPath}.corrupt-${now()}`; await rename(registryPath, backupPath).catch(() => undefined); return { sessions: {} }; } } async function save(state: TelegramSessionRegistryState): Promise { await mkdir(dirname(registryPath), { recursive: true }); await writeFile(registryPath, JSON.stringify({ sessions: state.sessions ?? {} }, null, 2) + "\n", { encoding: "utf8", mode: 0o600, }); } async function mutate( fn: (state: TelegramSessionRegistryState) => TelegramRegisteredSession | undefined | void, ): Promise { const state = await load(); const result = fn(state); await save(state); return result ?? undefined; } return { getPath: () => registryPath, getAgentDir: () => agentDir, load, save, register: async (session) => { const state = await load(); const existing = state.sessions[session.key]; const next: TelegramRegisteredSession = { ...existing, ...session, binding: session.binding ?? existing?.binding, inboxCursor: session.inboxCursor ?? existing?.inboxCursor, broadcastEnabled: session.broadcastEnabled ?? existing?.broadcastEnabled ?? false, updatedAt: session.updatedAt, }; state.sessions[session.key] = next; await mkdir(dirname(session.inboxPath), { recursive: true }); if (!existsSync(session.inboxPath)) await writeFile(session.inboxPath, "", { mode: 0o600 }); await save(state); return next; }, unregister: async (key) => { await mutate((state) => { delete state.sessions[key]; }); }, getSession: async (key) => (await load()).sessions[key], findByTopic: async (chatId, threadId) => { const state = await load(); return Object.values(state.sessions).find( (session) => session.binding?.chatId === chatId && session.binding?.threadId === threadId, ); }, setBindCode: async (key, code, ttlMs) => mutate((state) => { const session = state.sessions[key]; if (!session) return undefined; session.bindCode = normalizeBindCode(code); session.bindCodeExpiresAt = now() + ttlMs; session.updatedAt = now(); return session; }), bindTopicByCode: async (code, binding) => mutate((state) => { const normalized = normalizeBindCode(code); const session = Object.values(state.sessions).find( (candidate) => candidate.bindCode === normalized && (candidate.bindCodeExpiresAt ?? 0) >= now(), ); if (!session) return undefined; session.binding = binding; delete session.bindCode; delete session.bindCodeExpiresAt; session.updatedAt = now(); return session; }), bindTopic: async (key, binding) => mutate((state) => { const session = state.sessions[key]; if (!session) return undefined; session.binding = binding; session.updatedAt = now(); return session; }), unbindTopic: async (key) => mutate((state) => { const session = state.sessions[key]; if (!session) return undefined; delete session.binding; session.updatedAt = now(); return session; }), setBroadcast: async (key, enabled) => mutate((state) => { const session = state.sessions[key]; if (!session) return undefined; session.broadcastEnabled = enabled; session.updatedAt = now(); return session; }), updateInboxCursor: async (key, cursor) => mutate((state) => { const session = state.sessions[key]; if (!session) return undefined; session.inboxCursor = Math.max(0, cursor); session.updatedAt = now(); return session; }), }; }