threecx.py 9.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289
  1. import asyncio
  2. import json
  3. import logging
  4. import ssl
  5. import time
  6. from pathlib import Path
  7. from typing import Any, Awaitable, Callable
  8. import httpx
  9. import websockets
  10. from .config import Settings
  11. log = logging.getLogger(__name__)
  12. class ThreeCXClient:
  13. def __init__(self, settings: Settings):
  14. self.s = settings
  15. self._token = None
  16. self._expires_at = 0.0
  17. self._token_lock = asyncio.Lock()
  18. async def token(self) -> str:
  19. if self._token and time.time() < self._expires_at - 60:
  20. return self._token
  21. async with self._token_lock:
  22. if self._token and time.time() < self._expires_at - 60:
  23. return self._token
  24. secret = Path(self.s.threecx_api_key).read_text().strip()
  25. async with httpx.AsyncClient(
  26. verify=self.s.threecx_verify_tls,
  27. timeout=15,
  28. ) as client:
  29. response = await client.post(
  30. f"{self.s.threecx_base_url}/connect/token",
  31. data={
  32. "client_id": self.s.threecx_client_id,
  33. "client_secret": secret,
  34. "grant_type": "client_credentials",
  35. },
  36. )
  37. response.raise_for_status()
  38. data = response.json()
  39. self._token = data["access_token"]
  40. self._expires_at = time.time() + int(data.get("expires_in", 3600))
  41. return self._token
  42. async def invalidate_token(self) -> None:
  43. async with self._token_lock:
  44. self._token = None
  45. self._expires_at = 0.0
  46. async def get_dn(self, dn: str) -> dict[str, Any]:
  47. for attempt in range(2):
  48. token = await self.token()
  49. async with httpx.AsyncClient(
  50. verify=self.s.threecx_verify_tls,
  51. timeout=15,
  52. ) as client:
  53. response = await client.get(
  54. f"{self.s.threecx_base_url}/callcontrol/{dn}",
  55. headers={"Authorization": f"Bearer {token}"},
  56. )
  57. if response.status_code == 401 and attempt == 0:
  58. log.warning("3CX token rejected for get_dn(%s); refreshing token", dn)
  59. await self.invalidate_token()
  60. continue
  61. response.raise_for_status()
  62. return response.json()
  63. raise RuntimeError("3CX get_dn failed after token refresh")
  64. async def get_participants(self, dn: str) -> list[dict[str, Any]]:
  65. for attempt in range(2):
  66. token = await self.token()
  67. async with httpx.AsyncClient(
  68. verify=self.s.threecx_verify_tls,
  69. timeout=15,
  70. ) as client:
  71. response = await client.get(
  72. f"{self.s.threecx_base_url}/callcontrol/{dn}/participants",
  73. headers={"Authorization": f"Bearer {token}"},
  74. )
  75. if response.status_code == 401 and attempt == 0:
  76. log.warning(
  77. "3CX token rejected for get_participants(%s); refreshing token",
  78. dn,
  79. )
  80. await self.invalidate_token()
  81. continue
  82. response.raise_for_status()
  83. data = response.json()
  84. if isinstance(data, list):
  85. return data
  86. return data.get("participants", [])
  87. raise RuntimeError("3CX get_participants failed after token refresh")
  88. async def make_call(
  89. self,
  90. dn: str,
  91. device_id: str,
  92. destination: str,
  93. timeout_sec: int = 30,
  94. ) -> dict[str, Any]:
  95. for attempt in range(2):
  96. token = await self.token()
  97. async with httpx.AsyncClient(
  98. verify=self.s.threecx_verify_tls,
  99. timeout=15,
  100. ) as client:
  101. response = await client.post(
  102. f"{self.s.threecx_base_url}/callcontrol/{dn}/devices/{device_id}/makecall",
  103. headers={"Authorization": f"Bearer {token}"},
  104. json={
  105. "destination": destination,
  106. "timeoutSec": timeout_sec,
  107. },
  108. )
  109. if response.status_code == 401 and attempt == 0:
  110. log.warning(
  111. "3CX token rejected for make_call(%s -> %s); refreshing token",
  112. dn,
  113. destination,
  114. )
  115. await self.invalidate_token()
  116. continue
  117. response.raise_for_status()
  118. return response.json()
  119. raise RuntimeError("3CX make_call failed after token refresh")
  120. async def get_call_log(
  121. self,
  122. period_from: str,
  123. period_to: str,
  124. top: int = 500,
  125. skip: int = 0,
  126. ) -> dict[str, Any]:
  127. """Read current 3CX CDR/CallLog data via ReportCallLogData."""
  128. for attempt in range(2):
  129. token = await self.token()
  130. params = (
  131. f"periodFrom={period_from},"
  132. f"periodTo={period_to},"
  133. "sourceType=0,"
  134. "sourceFilter='',"
  135. "destinationType=0,"
  136. "destinationFilter='',"
  137. "callsType=0,"
  138. "callTimeFilterType=0,"
  139. "callTimeFilterFrom='0:00:0',"
  140. "callTimeFilterTo='0:00:0',"
  141. "hidePcalls=true"
  142. )
  143. url = (
  144. f"{self.s.threecx_base_url}/xapi/v1/"
  145. f"ReportCallLogData/Pbx.GetCallLogData({params})"
  146. f"?$top={top}&$skip={skip}&$count=true"
  147. )
  148. async with httpx.AsyncClient(
  149. verify=self.s.threecx_verify_tls,
  150. timeout=120,
  151. ) as client:
  152. response = await client.get(
  153. url,
  154. headers={
  155. "Authorization": f"Bearer {token}",
  156. "Accept": "application/json",
  157. },
  158. )
  159. if response.status_code == 401 and attempt == 0:
  160. log.warning("3CX token rejected for get_call_log; refreshing token")
  161. await self.invalidate_token()
  162. continue
  163. response.raise_for_status()
  164. return response.json()
  165. raise RuntimeError("3CX get_call_log failed after token refresh")
  166. async def download_recording(self, rec_id: int) -> tuple[bytes, str]:
  167. """Download a historical 3CX recording by recording ID."""
  168. token = await self.token()
  169. url = (
  170. f"{self.s.threecx_base_url}/xapi/v1/"
  171. f"Recordings/Pbx.DownloadRecording(recId={int(rec_id)})"
  172. )
  173. async with httpx.AsyncClient(
  174. verify=self.s.threecx_verify_tls,
  175. timeout=120,
  176. ) as client:
  177. response = await client.get(
  178. url,
  179. headers={
  180. "Authorization": f"Bearer {token}",
  181. "Accept": "audio/x-wav,*/*",
  182. },
  183. )
  184. if response.status_code == 401:
  185. await self.invalidate_token()
  186. token = await self.token()
  187. async with httpx.AsyncClient(
  188. verify=self.s.threecx_verify_tls,
  189. timeout=120,
  190. ) as client:
  191. response = await client.get(
  192. url,
  193. headers={
  194. "Authorization": f"Bearer {token}",
  195. "Accept": "audio/x-wav,*/*",
  196. },
  197. )
  198. response.raise_for_status()
  199. content_type = response.headers.get(
  200. "content-type",
  201. "audio/x-wav",
  202. )
  203. return response.content, content_type
  204. async def websocket(self, on_event: Callable[[dict[str, Any]], Awaitable[None]]):
  205. token = await self.token()
  206. uri = self.s.threecx_base_url.replace("https://", "wss://").replace("http://", "ws://")
  207. uri += "/callcontrol/ws"
  208. ssl_context = None
  209. if uri.startswith("wss://") and not self.s.threecx_verify_tls:
  210. ssl_context = ssl.create_default_context()
  211. ssl_context.check_hostname = False
  212. ssl_context.verify_mode = ssl.CERT_NONE
  213. async with websockets.connect(
  214. uri,
  215. additional_headers={"Authorization": f"Bearer {token}"},
  216. ssl=ssl_context,
  217. ping_interval=20,
  218. ping_timeout=20,
  219. ) as ws:
  220. await ws.send(json.dumps({
  221. "RequestID": "telephony-middleware",
  222. "Path": "/callcontrol",
  223. }))
  224. log.info("3CX WebSocket connected")
  225. async for raw in ws:
  226. try:
  227. message = json.loads(raw)
  228. except json.JSONDecodeError:
  229. continue
  230. await on_event(message)
  231. async def run_websocket(self, on_event):
  232. delay = 2
  233. while True:
  234. try:
  235. await self.websocket(on_event)
  236. delay = 2
  237. except asyncio.CancelledError:
  238. raise
  239. except Exception:
  240. log.exception("3CX WebSocket failed; reconnect in %ss", delay)
  241. await asyncio.sleep(delay)
  242. delay = min(delay * 2, 60)