main.py 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684
  1. from datetime import datetime, timezone
  2. import hashlib
  3. import httpx
  4. import logging
  5. import secrets
  6. import json
  7. from contextlib import asynccontextmanager
  8. from fastapi import FastAPI, Header, HTTPException, Query
  9. from fastapi.middleware.cors import CORSMiddleware
  10. from fastapi.responses import JSONResponse
  11. from pydantic import BaseModel
  12. from fastapi import Query
  13. from .phone import normalize
  14. from .calls import normalize_call
  15. from .config import Settings
  16. from .service import TelephonyService
  17. settings = Settings()
  18. service = TelephonyService(settings)
  19. class ResolveRequest(BaseModel):
  20. until: str
  21. by: str
  22. note: str | None = None
  23. class ReopenRequest(BaseModel):
  24. by: str
  25. def problem(status, title, detail, field=None, retryable=False):
  26. body = {
  27. "type": "about:blank",
  28. "title": title,
  29. "status": status,
  30. "detail": detail,
  31. "retryable": retryable,
  32. }
  33. if field:
  34. body["field"] = field
  35. return JSONResponse(
  36. status_code=status,
  37. content=body,
  38. media_type="application/problem+json",
  39. )
  40. def authorized(authorization):
  41. if not authorization:
  42. return False
  43. scheme, _, token = authorization.partition(" ")
  44. return (
  45. scheme.lower() == "bearer"
  46. and secrets.compare_digest(
  47. token,
  48. settings.middleware_api_token,
  49. )
  50. )
  51. @asynccontextmanager
  52. async def lifespan(app):
  53. logging.basicConfig(
  54. level=getattr(
  55. logging,
  56. settings.log_level.upper(),
  57. logging.INFO,
  58. ),
  59. format="%(asctime)s %(levelname)s %(name)s: %(message)s",
  60. )
  61. await service.start()
  62. yield
  63. class OutboundCallRequest(BaseModel):
  64. sourceDn: str
  65. destination: str
  66. timeoutSec: int = 30
  67. app = FastAPI(
  68. title="3CX Telefonie Middleware",
  69. version="0.2.0",
  70. lifespan=lifespan,
  71. )
  72. app.add_middleware(
  73. CORSMiddleware,
  74. allow_origins=[
  75. x.strip()
  76. for x in settings.cors_origins.split(",")
  77. if x.strip()
  78. ],
  79. allow_credentials=False,
  80. allow_methods=["GET", "POST", "OPTIONS"],
  81. allow_headers=[
  82. "Authorization",
  83. "Content-Type",
  84. "Idempotency-Key",
  85. ],
  86. )
  87. @app.get("/phone/normalize")
  88. async def phone_normalize(
  89. number: str = Query(..., min_length=1),
  90. country: str = Query("DE", min_length=2, max_length=2),
  91. ):
  92. return normalize(number, country.upper())
  93. @app.get("/calls/active")
  94. async def active_calls():
  95. queue_dn = settings.threecx_queue_dn
  96. participants = service.get_active_participants()
  97. data = [
  98. normalize_call(participant, queue_dn)
  99. for participant in participants
  100. if participant.get("party_dn_type") == "Wexternalline"
  101. ]
  102. return {
  103. "data": data,
  104. "meta": {
  105. "totalItems": len(data)
  106. }
  107. }
  108. @app.post("/calls/outbound")
  109. async def outbound_call(
  110. request: OutboundCallRequest,
  111. authorization: str | None = Header(default=None),
  112. idempotency_key: str | None = Header(default=None, alias="Idempotency-Key"),
  113. ):
  114. if not authorized(authorization):
  115. return problem(
  116. 401,
  117. "Unauthorized",
  118. "Bearer-Token fehlt oder ist ungültig.",
  119. retryable=False,
  120. )
  121. if not idempotency_key or not idempotency_key.strip():
  122. return problem(
  123. 400,
  124. "Idempotency-Key fehlt",
  125. "Für ausgehende Anrufe ist ein Idempotency-Key erforderlich.",
  126. retryable=False,
  127. )
  128. idempotency_key = idempotency_key.strip()
  129. if len(idempotency_key) > 200:
  130. return problem(
  131. 400,
  132. "Idempotency-Key ungültig",
  133. "Der Idempotency-Key darf maximal 200 Zeichen lang sein.",
  134. retryable=False,
  135. )
  136. try:
  137. dn = request.sourceDn.strip()
  138. if not dn.isdigit():
  139. return problem(
  140. 422,
  141. "Ungültige Nebenstelle",
  142. "sourceDn muss eine numerische 3CX-Nebenstelle sein.",
  143. retryable=False,
  144. )
  145. if dn not in settings.monitored_dns:
  146. return problem(
  147. 403,
  148. "Nebenstelle nicht freigegeben",
  149. f"Die Nebenstelle {dn} ist für ausgehende API-Anrufe nicht freigegeben.",
  150. retryable=False,
  151. )
  152. raw_destination = request.destination.strip()
  153. # 1-4-stellige numerische Ziele sind interne 3CX-Nebenstellen.
  154. if raw_destination.isdigit() and 1 <= len(raw_destination) <= 4:
  155. normalized = {
  156. "raw": raw_destination,
  157. "e164": raw_destination,
  158. "valid": True,
  159. "type": "EXTENSION",
  160. }
  161. else:
  162. normalized = normalize(
  163. raw_destination,
  164. settings.phone_default_region,
  165. )
  166. if not normalized.get("valid"):
  167. return problem(
  168. 422,
  169. "Ungültige Rufnummer",
  170. "Die Zielrufnummer konnte nicht gültig normalisiert werden.",
  171. retryable=False,
  172. )
  173. request_fingerprint = hashlib.sha256(
  174. json.dumps(
  175. {
  176. "sourceDn": dn,
  177. "destination": normalized.get("e164"),
  178. "timeoutSec": request.timeoutSec,
  179. },
  180. sort_keys=True,
  181. separators=(",", ":"),
  182. ).encode("utf-8")
  183. ).hexdigest()
  184. existing = await service.repo.get_outbound_idempotency(idempotency_key)
  185. if existing:
  186. if existing["request_hash"] != request_fingerprint:
  187. return problem(
  188. 409,
  189. "Idempotency-Key bereits verwendet",
  190. "Der Idempotency-Key wurde bereits für eine andere Anfrage verwendet.",
  191. retryable=False,
  192. )
  193. if existing["status"] == "completed" and existing["response_json"]:
  194. return json.loads(existing["response_json"])
  195. return problem(
  196. 409,
  197. "Anruf wird bereits verarbeitet",
  198. "Für diesen Idempotency-Key läuft bereits ein Anrufvorgang.",
  199. retryable=True,
  200. )
  201. claimed = await service.repo.claim_outbound_idempotency(
  202. idempotency_key,
  203. request_fingerprint,
  204. datetime.now(timezone.utc).isoformat(),
  205. )
  206. if not claimed:
  207. existing = await service.repo.get_outbound_idempotency(idempotency_key)
  208. if existing and existing["request_hash"] != request_fingerprint:
  209. return problem(
  210. 409,
  211. "Idempotency-Key bereits verwendet",
  212. "Der Idempotency-Key wurde bereits für eine andere Anfrage verwendet.",
  213. retryable=False,
  214. )
  215. return problem(
  216. 409,
  217. "Anruf wird bereits verarbeitet",
  218. "Für diesen Idempotency-Key läuft bereits ein Anrufvorgang.",
  219. retryable=True,
  220. )
  221. dn_data = await service.cx.get_dn(dn)
  222. devices = dn_data.get("devices") or []
  223. if len(devices) == 0:
  224. return problem(
  225. 409,
  226. "Kein Gerät gefunden",
  227. f"Für Nebenstelle {dn} wurde kein 3CX-Gerät gefunden.",
  228. retryable=True,
  229. )
  230. hint = settings.outbound_device_hints.get(dn)
  231. if not hint:
  232. return problem(
  233. 409,
  234. "Kein Hardware-Gerät konfiguriert",
  235. f"Für Nebenstelle {dn} ist kein Hardware-Telefon für ausgehende API-Anrufe konfiguriert.",
  236. retryable=False,
  237. )
  238. matching_devices = [
  239. device
  240. for device in devices
  241. if hint.lower() in str(device.get("user_agent", "")).lower()
  242. ]
  243. if len(matching_devices) == 0:
  244. return problem(
  245. 409,
  246. "Hardware-Gerät nicht gefunden",
  247. f"Für Nebenstelle {dn} wurde kein 3CX-Gerät passend zu '{hint}' gefunden.",
  248. retryable=True,
  249. )
  250. if len(matching_devices) > 1:
  251. return problem(
  252. 409,
  253. "Hardware-Gerät nicht eindeutig",
  254. f"Für Nebenstelle {dn} wurden mehrere Geräte passend zu '{hint}' gefunden.",
  255. retryable=False,
  256. )
  257. device = matching_devices[0]
  258. device_id = (
  259. device.get("device_id")
  260. or device.get("deviceId")
  261. or device.get("id")
  262. )
  263. if not device_id:
  264. return problem(
  265. 502,
  266. "Ungültige 3CX-Geräteantwort",
  267. "3CX hat kein verwendbares Gerätekennzeichen geliefert.",
  268. retryable=True,
  269. )
  270. result = await service.cx.make_call(
  271. dn=dn,
  272. device_id=str(device_id),
  273. destination=normalized["e164"],
  274. timeout_sec=request.timeoutSec,
  275. )
  276. response = {
  277. "status": "accepted",
  278. "sourceDn": dn,
  279. "deviceId": str(device_id),
  280. "destination": {
  281. "raw": normalized.get("raw"),
  282. "e164": normalized.get("e164"),
  283. "valid": normalized.get("valid"),
  284. "type": normalized.get("type"),
  285. },
  286. "threecx": result,
  287. }
  288. threecx_result = result.get("result") or result
  289. threecx_callid = threecx_result.get("callid")
  290. threecx_legid = threecx_result.get("legid")
  291. if threecx_callid is not None:
  292. started_at = datetime.now(timezone.utc).isoformat()
  293. await service.repo.insert_outbound_call(
  294. callid=threecx_callid,
  295. legid=threecx_legid,
  296. source_dn=dn,
  297. destination=normalized.get("e164"),
  298. phone=(
  299. normalized.get("e164")
  300. if normalized.get("type") != "EXTENSION"
  301. else None
  302. ),
  303. raw=normalized.get("raw"),
  304. started_at=started_at,
  305. )
  306. await service.repo.finish_outbound_idempotency(
  307. idempotency_key,
  308. json.dumps(response, ensure_ascii=False, separators=(",", ":")),
  309. )
  310. return response
  311. except httpx.HTTPStatusError as exc:
  312. await service.repo.release_outbound_idempotency(idempotency_key)
  313. status = exc.response.status_code
  314. if status in (401, 403):
  315. return problem(
  316. 502,
  317. "3CX-Authentifizierung fehlgeschlagen",
  318. "Die Middleware konnte den 3CX-Anruf nicht authentifizieren.",
  319. retryable=False,
  320. )
  321. if status == 404:
  322. return problem(
  323. 502,
  324. "3CX-Gerät oder Ziel nicht gefunden",
  325. "3CX hat die angeforderte Ressource nicht gefunden.",
  326. retryable=False,
  327. )
  328. return problem(
  329. 502,
  330. "3CX-Anruf fehlgeschlagen",
  331. f"3CX hat HTTP {status} zurückgegeben.",
  332. retryable=True,
  333. )
  334. except Exception:
  335. await service.repo.release_outbound_idempotency(idempotency_key)
  336. log.exception("Outbound call failed")
  337. return problem(
  338. 502,
  339. "Ausgehender Anruf fehlgeschlagen",
  340. "Der Anruf konnte über 3CX nicht ausgelöst werden.",
  341. retryable=True,
  342. )
  343. @app.get("/calls/cdr")
  344. async def calls_cdr(
  345. from_: str = Query(..., alias="from"),
  346. to: str = Query(...),
  347. limit: int = Query(500, ge=1, le=1000),
  348. offset: int = Query(0, ge=0),
  349. authorization: str | None = Header(default=None),
  350. ):
  351. if not authorized(authorization):
  352. return problem(
  353. 401,
  354. "Unauthorized",
  355. "Bearer-Token fehlt oder ist ungültig.",
  356. retryable=False,
  357. )
  358. try:
  359. data = await service.cx.get_call_log(
  360. period_from=from_,
  361. period_to=to,
  362. top=limit,
  363. skip=offset,
  364. )
  365. return {
  366. "data": data.get("value", []),
  367. "meta": {
  368. "totalItems": data.get("@odata.count"),
  369. "limit": limit,
  370. "offset": offset,
  371. },
  372. }
  373. except Exception:
  374. log.exception("3CX CDR query failed")
  375. return problem(
  376. 502,
  377. "3CX CDR-Abfrage fehlgeschlagen",
  378. "Die Call-History konnte von 3CX nicht abgerufen werden.",
  379. retryable=True,
  380. )
  381. @app.get("/calls/recording/{rec_id}")
  382. async def calls_recording(
  383. rec_id: int,
  384. authorization: str | None = Header(default=None),
  385. ):
  386. if not authorized(authorization):
  387. return problem(
  388. 401,
  389. "Unauthorized",
  390. "Bearer-Token fehlt oder ist ungültig.",
  391. retryable=False,
  392. )
  393. try:
  394. from fastapi.responses import Response
  395. content, content_type = await service.cx.download_recording(rec_id)
  396. return Response(
  397. content=content,
  398. media_type=content_type.split(";", 1)[0],
  399. headers={
  400. "Content-Disposition": (
  401. f'attachment; filename="recording-{rec_id}.wav"'
  402. )
  403. },
  404. )
  405. except Exception:
  406. log.exception(
  407. "3CX recording download failed: rec_id=%s",
  408. rec_id,
  409. )
  410. return problem(
  411. 502,
  412. "Recording konnte nicht geladen werden",
  413. "Das 3CX-Recording konnte nicht abgerufen werden.",
  414. retryable=True,
  415. )
  416. @app.get("/calls/{call_id}")
  417. async def call_detail(call_id: int):
  418. row = await service.repo.get_call(call_id)
  419. if row is None:
  420. return problem(
  421. 404,
  422. "Call nicht gefunden",
  423. f"Call {call_id} wurde nicht gefunden.",
  424. retryable=False,
  425. )
  426. return row
  427. @app.get("/calls")
  428. async def calls_history(
  429. limit: int = Query(100, ge=1, le=500),
  430. offset: int = Query(0, ge=0),
  431. status: str | None = Query(None),
  432. ):
  433. rows, total = await service.repo.list_calls(
  434. limit=limit,
  435. offset=offset,
  436. status=status,
  437. )
  438. return {
  439. "data": rows,
  440. "meta": {
  441. "totalItems": total,
  442. "limit": limit,
  443. "offset": offset,
  444. }
  445. }
  446. @app.get("/health")
  447. async def health():
  448. return {
  449. "status": "ok",
  450. "queue": settings.threecx_queue_dn,
  451. "version": "0.2.0",
  452. }
  453. @app.get("/missed-calls")
  454. async def missed_calls():
  455. return await service.repo.list_groups()
  456. @app.get("/missed-calls/groups")
  457. async def missed_call_groups():
  458. groups = await service.repo.list_groups(open_only=True)
  459. return {
  460. "data": groups,
  461. "meta": {
  462. "totalItems": len(groups),
  463. },
  464. }
  465. @app.get("/missed-calls/groups/{e164}")
  466. async def missed_call_group(e164: str):
  467. group = await service.repo.get_group(e164)
  468. if not group:
  469. return problem(
  470. 404,
  471. "Not found",
  472. "Die Rufnummer wurde nicht gefunden.",
  473. )
  474. return group
  475. @app.post("/missed-calls/groups/{e164}/resolve")
  476. async def resolve_group(
  477. e164: str,
  478. payload: ResolveRequest,
  479. authorization: str | None = Header(default=None),
  480. idempotency_key: str | None = Header(default=None),
  481. ):
  482. if not authorized(authorization):
  483. return problem(
  484. 401,
  485. "Unauthorized",
  486. "Bearer-Token fehlt oder ist ungültig.",
  487. )
  488. if not idempotency_key:
  489. return problem(
  490. 400,
  491. "Idempotency-Key required",
  492. "Idempotency-Key ist erforderlich.",
  493. field="Idempotency-Key",
  494. )
  495. previous = await service.repo.get_idempotency(idempotency_key, "resolve")
  496. if previous:
  497. return JSONResponse(
  498. content=json.loads(previous),
  499. status_code=200,
  500. )
  501. group = await service.repo.get_group(e164)
  502. if not group:
  503. return problem(
  504. 404,
  505. "Not found",
  506. "Die Rufnummer wurde nicht gefunden.",
  507. )
  508. result = await service.repo.resolve(
  509. e164,
  510. payload.until,
  511. payload.by,
  512. payload.note,
  513. )
  514. response_json = json.dumps(result)
  515. await service.repo.save_idempotency(
  516. idempotency_key,
  517. "resolve",
  518. response_json,
  519. )
  520. return result
  521. @app.post("/missed-calls/groups/{e164}/reopen")
  522. async def reopen_group(
  523. e164: str,
  524. payload: ReopenRequest,
  525. authorization: str | None = Header(default=None),
  526. idempotency_key: str | None = Header(default=None),
  527. ):
  528. if not authorized(authorization):
  529. return problem(
  530. 401,
  531. "Unauthorized",
  532. "Bearer-Token fehlt oder ist ungültig.",
  533. )
  534. if not idempotency_key:
  535. return problem(
  536. 400,
  537. "Idempotency-Key required",
  538. "Idempotency-Key ist erforderlich.",
  539. field="Idempotency-Key",
  540. )
  541. previous = await service.repo.get_idempotency(
  542. idempotency_key,
  543. "reopen",
  544. )
  545. if previous:
  546. return JSONResponse(
  547. content=json.loads(previous),
  548. status_code=200,
  549. )
  550. group = await service.repo.get_group(e164)
  551. if not group:
  552. return problem(
  553. 404,
  554. "Not found",
  555. "Die Rufnummer wurde nicht gefunden.",
  556. )
  557. result = await service.repo.reopen(
  558. e164,
  559. payload.by,
  560. )
  561. response_json = json.dumps(result)
  562. await service.repo.save_idempotency(
  563. idempotency_key,
  564. "reopen",
  565. response_json,
  566. )
  567. return result