/** * Telegram singleton lock helpers * Owns shared locks.json access and Telegram bridge ownership semantics */ import { existsSync, mkdirSync, readFileSync, renameSync, unlinkSync, writeFileSync, } from "node:fs"; import { homedir } from "node:os"; import { dirname, join, resolve } from "node:path"; export const TELEGRAM_LOCK_KEY = "pi-telegram-group-topic"; export const TELEGRAM_LEGACY_LOCK_KEYS = ["@llblab/pi-telegram"]; function getAgentDir(): string { return process.env.PI_CODING_AGENT_DIR ? resolve(process.env.PI_CODING_AGENT_DIR) : join(homedir(), ".pi", "agent"); } function getLocksPath(): string { return join(getAgentDir(), "locks.json"); } export interface TelegramLockEntry { pid: number; cwd?: string; } export interface TelegramLockContext { cwd: string; } export type TelegramLockState = | { kind: "inactive" } | { kind: "active-here"; lock: TelegramLockEntry } | { kind: "active-elsewhere"; lock: TelegramLockEntry } | { kind: "stale"; lock: TelegramLockEntry }; export interface TelegramLockAcquireOptions { force?: boolean; } export type TelegramLockAcquireResult = | { ok: true; lock: TelegramLockEntry; replacedStale: boolean } | { ok: false; lock: TelegramLockEntry }; export interface TelegramLockRuntime { acquire: ( ctx: TContext, options?: TelegramLockAcquireOptions, ) => TelegramLockAcquireResult; release: () => TelegramLockState; getState: () => TelegramLockState; getStatusLabel: () => string; owns: (ctx?: TelegramLockContext) => boolean; } export interface TelegramLockRuntimeOptions { key?: string; legacyKeys?: string[]; locksPath?: string; pid?: number; isProcessAlive?: (pid: number) => boolean; } export interface TelegramLockedPollingStartOptions { force?: boolean; } export type TelegramLockedPollingStartResult = | { ok: true; message: string; canTakeover?: false } | { ok: false; message: string; canTakeover?: boolean; owner?: string }; export interface TelegramLockedPollingRuntime< TContext extends TelegramLockContext, > { start: ( ctx: TContext, options?: TelegramLockedPollingStartOptions, ) => Promise; stop: () => Promise; suspend: () => Promise; onSessionStart: (_event: unknown, ctx: TContext) => Promise; } export interface TelegramLockedPollingRuntimeDeps< TContext extends TelegramLockContext, > { lock: TelegramLockRuntime; hasBotToken: () => boolean; startPolling: (ctx: TContext) => void | Promise; stopPolling: () => Promise; updateStatus: (ctx: TContext) => void; recordRuntimeEvent?: ( category: string, error: unknown, details?: Record, ) => void; ownershipCheckMs?: number; } export function readLocks(path = getLocksPath()): Record { if (!existsSync(path)) return {}; try { const value = JSON.parse(readFileSync(path, "utf8")); return value && typeof value === "object" && !Array.isArray(value) ? (value as Record) : {}; } catch { return {}; } } export function writeLocks(path: string, locks: Record): void { mkdirSync(dirname(path), { recursive: true }); const tempPath = `${path}.${process.pid}.${Date.now()}.tmp`; try { writeFileSync(tempPath, `${JSON.stringify(locks, null, 2)}\n`, "utf8"); renameSync(tempPath, path); } catch (error) { try { unlinkSync(tempPath); } catch { /* best effort */ } throw error; } } export function parseTelegramLockEntry( value: unknown, ): TelegramLockEntry | undefined { if (!value || typeof value !== "object" || Array.isArray(value)) return undefined; const record = value as Record; if (typeof record.pid !== "number") return undefined; return { pid: record.pid, cwd: typeof record.cwd === "string" ? record.cwd : undefined, }; } export function isProcessAlive(pid: number): boolean { if (!Number.isInteger(pid) || pid <= 0) return false; try { process.kill(pid, 0); return true; } catch (error) { return (error as { code?: string }).code === "EPERM"; } } function formatLock(lock: TelegramLockEntry): string { return lock.cwd ? `pid ${lock.pid}, cwd ${lock.cwd}` : `pid ${lock.pid}`; } function getLockState( lock: TelegramLockEntry | undefined, pid: number, isAlive: (pid: number) => boolean, ): TelegramLockState { return getCombinedLockState(lock ? [lock] : [], pid, isAlive); } function getCombinedLockState( locks: TelegramLockEntry[], pid: number, isAlive: (pid: number) => boolean, ): TelegramLockState { if (locks.length === 0) return { kind: "inactive" }; const liveElsewhere = locks.find((lock) => lock.pid !== pid && isAlive(lock.pid)); if (liveElsewhere) return { kind: "active-elsewhere", lock: liveElsewhere }; const activeHere = locks.find((lock) => lock.pid === pid); if (activeHere) return { kind: "active-here", lock: activeHere }; return { kind: "stale", lock: locks[0] }; } function ownsLockContext( lock: TelegramLockEntry | undefined, pid: number, ctx?: TelegramLockContext, ): boolean { if (!lock || lock.pid !== pid) return false; return !lock.cwd || !ctx || lock.cwd === ctx.cwd; } function snapshotLockContext(ctx: TelegramLockContext): TelegramLockContext { return { cwd: ctx.cwd }; } function formatLockState(state: TelegramLockState): string { switch (state.kind) { case "inactive": return "inactive"; case "active-here": return "active here"; case "active-elsewhere": return `active elsewhere (${formatLock(state.lock)})`; case "stale": return `stale (${formatLock(state.lock)})`; } } export function createTelegramLockRuntime( options: TelegramLockRuntimeOptions = {}, ): TelegramLockRuntime { const key = options.key ?? TELEGRAM_LOCK_KEY; const legacyKeys = options.legacyKeys ?? (options.key ? [] : TELEGRAM_LEGACY_LOCK_KEYS); const lockKeys = [key, ...legacyKeys.filter((legacyKey) => legacyKey !== key)]; const locksPath = options.locksPath ?? getLocksPath(); const pid = options.pid ?? process.pid; const isAlive = options.isProcessAlive ?? isProcessAlive; const readState = () => { const locks = readLocks(locksPath); return getCombinedLockState( lockKeys .map((lockKey) => parseTelegramLockEntry(locks[lockKey])) .filter((lock): lock is TelegramLockEntry => !!lock), pid, isAlive, ); }; const writeLock = (lock: TelegramLockEntry) => { const locks = readLocks(locksPath); for (const lockKey of lockKeys) locks[lockKey] = lock; writeLocks(locksPath, locks); }; return { acquire: (ctx, acquireOptions = {}) => { const state = readState(); if (state.kind === "active-elsewhere" && !acquireOptions.force) return { ok: false, lock: state.lock }; const lock = { pid, cwd: ctx.cwd }; writeLock(lock); return { ok: true, lock, replacedStale: state.kind === "stale" }; }, release: () => { const state = readState(); if (state.kind === "active-here" || state.kind === "stale") { const locks = readLocks(locksPath); for (const lockKey of lockKeys) delete locks[lockKey]; writeLocks(locksPath, locks); } return state; }, getState: readState, getStatusLabel: () => formatLockState(readState()), owns: (ctx) => { const state = readState(); return ownsLockContext(state.kind === "active-here" ? state.lock : undefined, pid, ctx); }, }; } export function createTelegramLockedPollingRuntime< TContext extends TelegramLockContext, >( deps: TelegramLockedPollingRuntimeDeps, ): TelegramLockedPollingRuntime { let ownershipInterval: ReturnType | undefined; let ownershipStop: Promise | undefined; const ownershipCheckMs = deps.ownershipCheckMs ?? 1000; const stopOwnershipWatcher = () => { if (!ownershipInterval) return; clearInterval(ownershipInterval); ownershipInterval = undefined; }; const suspendPolling = async () => { stopOwnershipWatcher(); if (ownershipStop) { await ownershipStop; return; } await deps.stopPolling(); }; const stopAfterOwnershipLoss = () => { if (ownershipStop) return; stopOwnershipWatcher(); ownershipStop = deps .stopPolling() .catch((error) => deps.recordRuntimeEvent?.("lock", error, { phase: "ownership-loss" }), ) .finally(() => { ownershipStop = undefined; }); }; const startOwnershipWatcher = (ctx: TContext) => { const owner = snapshotLockContext(ctx); stopOwnershipWatcher(); ownershipInterval = setInterval(() => { if (deps.lock.owns(owner)) return; stopAfterOwnershipLoss(); }, ownershipCheckMs); ownershipInterval.unref?.(); }; return { start: async (ctx, options = {}) => { if (!deps.hasBotToken()) return { ok: false, message: "Telegram bot is not configured." }; const acquired = deps.lock.acquire(ctx, options); if (!acquired.ok) { return { ok: false, canTakeover: true, owner: formatLock(acquired.lock), message: `Telegram bridge is active in another pi instance (${formatLock(acquired.lock)}).`, }; } await deps.startPolling(ctx); startOwnershipWatcher(ctx); deps.updateStatus(ctx); const staleSuffix = acquired.replacedStale ? " Replaced stale lock." : ""; return { ok: true, message: `Telegram bridge connected.${staleSuffix}` }; }, stop: async () => { await suspendPolling(); const state = deps.lock.release(); if (state.kind === "active-elsewhere") { return `Telegram bridge is active in another pi instance (${formatLock(state.lock)}).`; } if (state.kind === "stale") return `Removed stale Telegram bridge lock (${formatLock(state.lock)}).`; return "Telegram bridge disconnected."; }, suspend: suspendPolling, onSessionStart: async (_event, ctx) => { if (!deps.hasBotToken()) return; const ownsCurrentLock = deps.lock.owns(ctx); const state = ownsCurrentLock ? undefined : deps.lock.getState(); const canResumeStaleSameCwd = state?.kind === "stale" && state.lock.cwd === ctx.cwd; if (!ownsCurrentLock && !canResumeStaleSameCwd) return; try { if (canResumeStaleSameCwd) { const acquired = deps.lock.acquire(ctx); if (!acquired.ok) return; } await deps.startPolling(ctx); startOwnershipWatcher(ctx); deps.updateStatus(ctx); } catch (error) { deps.recordRuntimeEvent?.("lock", error, { phase: "auto-start" }); } }, }; }