111 lines
4.1 KiB
Python
111 lines
4.1 KiB
Python
import asyncio
|
|
import logging
|
|
import os
|
|
|
|
from pyrogram import Client
|
|
from pytgcalls import PyTgCalls, filters as pfl
|
|
from pytgcalls.exceptions import CallBusy, CallDeclined, CallDiscarded, TimedOutAnswer
|
|
from pytgcalls.types import CallConfig, MediaStream, StreamEnded
|
|
|
|
from app.config import settings
|
|
from app.tts import synthesize
|
|
|
|
logger = logging.getLogger("callsvc.calling")
|
|
|
|
pyro_client = Client(
|
|
settings.session_name,
|
|
api_id=settings.api_id,
|
|
api_hash=settings.api_hash,
|
|
phone_number=settings.phone_number,
|
|
workdir=settings.sessions_dir,
|
|
)
|
|
|
|
call_py = PyTgCalls(pyro_client)
|
|
|
|
# per-цель "stream finished" события, заполняются on_stream_end хендлером
|
|
_stream_finished: dict[int, asyncio.Event] = {}
|
|
# один вызов на цель одновременно, чтобы не пересекались звонки
|
|
_locks: dict[str, asyncio.Lock] = {}
|
|
|
|
|
|
def _lock_for(key: str) -> asyncio.Lock:
|
|
if key not in _locks:
|
|
_locks[key] = asyncio.Lock()
|
|
return _locks[key]
|
|
|
|
|
|
@call_py.on_update(pfl.stream_end)
|
|
async def _on_stream_end(_: PyTgCalls, update: StreamEnded):
|
|
event = _stream_finished.get(update.chat_id)
|
|
if event is not None:
|
|
event.set()
|
|
|
|
|
|
async def start():
|
|
await pyro_client.start()
|
|
await call_py.start()
|
|
logger.info("Telegram user client + PyTgCalls started")
|
|
|
|
|
|
async def stop():
|
|
# PyTgCalls не имеет отдельного stop() — он поднимает/использует тот же
|
|
# pyrogram-клиент, поэтому останавливаем только клиент.
|
|
await pyro_client.stop()
|
|
logger.info("Telegram user client + PyTgCalls stopped")
|
|
|
|
|
|
async def make_call(text: str, target: int | str | None = None) -> dict:
|
|
"""Звонит пользователю настоящим p2p-звонком (с гудком и ожиданием ответа)
|
|
и озвучивает текст (TTS) сразу после того, как собеседник ответил.
|
|
"""
|
|
target = target if target is not None else settings.call_target
|
|
lock = _lock_for(str(target))
|
|
|
|
if lock.locked():
|
|
raise RuntimeError(f"Call to {target} is already in progress")
|
|
|
|
async with lock:
|
|
audio_path, duration = synthesize(text)
|
|
resolved_id = await call_py.resolve_chat_id(target)
|
|
if resolved_id <= 0:
|
|
raise RuntimeError(
|
|
f"{target} resolves to a group/channel id ({resolved_id}). "
|
|
"This service only makes direct p2p calls to users."
|
|
)
|
|
event = asyncio.Event()
|
|
_stream_finished[resolved_id] = event
|
|
|
|
try:
|
|
config = CallConfig(timeout=settings.call_ring_timeout)
|
|
try:
|
|
logger.info("Calling %s (resolved id %s)...", target, resolved_id)
|
|
await call_py.play(resolved_id, MediaStream(audio_path), config=config)
|
|
except TimedOutAnswer as exc:
|
|
raise RuntimeError(f"{target} did not answer the call in time") from exc
|
|
except CallDeclined as exc:
|
|
raise RuntimeError(f"{target} declined the call") from exc
|
|
except CallBusy as exc:
|
|
raise RuntimeError(f"{target} is busy") from exc
|
|
except CallDiscarded as exc:
|
|
raise RuntimeError(f"Call to {target} was discarded") from exc
|
|
|
|
logger.info("Call answered, playing TTS (%.1fs) to %s", duration, target)
|
|
|
|
# ждём сигнал о завершении стрима, но не дольше duration + запас
|
|
timeout = duration + 10
|
|
try:
|
|
await asyncio.wait_for(event.wait(), timeout=timeout)
|
|
except asyncio.TimeoutError:
|
|
logger.warning(
|
|
"Stream end event not received for %s within %.1fs, hanging up anyway",
|
|
target,
|
|
timeout,
|
|
)
|
|
|
|
await call_py.leave_call(resolved_id)
|
|
return {"target": target, "duration": duration}
|
|
finally:
|
|
_stream_finished.pop(resolved_id, None)
|
|
if os.path.exists(audio_path):
|
|
os.remove(audio_path)
|