| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684 |
- from datetime import datetime, timezone
- import hashlib
- import httpx
- import logging
- import secrets
- import json
- from contextlib import asynccontextmanager
- from fastapi import FastAPI, Header, HTTPException, Query
- from fastapi.middleware.cors import CORSMiddleware
- from fastapi.responses import JSONResponse
- from pydantic import BaseModel
- from fastapi import Query
- from .phone import normalize
- from .calls import normalize_call
- from .config import Settings
- from .service import TelephonyService
- settings = Settings()
- service = TelephonyService(settings)
- class ResolveRequest(BaseModel):
- until: str
- by: str
- note: str | None = None
- class ReopenRequest(BaseModel):
- by: str
- def problem(status, title, detail, field=None, retryable=False):
- body = {
- "type": "about:blank",
- "title": title,
- "status": status,
- "detail": detail,
- "retryable": retryable,
- }
- if field:
- body["field"] = field
- return JSONResponse(
- status_code=status,
- content=body,
- media_type="application/problem+json",
- )
- def authorized(authorization):
- if not authorization:
- return False
- scheme, _, token = authorization.partition(" ")
- return (
- scheme.lower() == "bearer"
- and secrets.compare_digest(
- token,
- settings.middleware_api_token,
- )
- )
- @asynccontextmanager
- async def lifespan(app):
- logging.basicConfig(
- level=getattr(
- logging,
- settings.log_level.upper(),
- logging.INFO,
- ),
- format="%(asctime)s %(levelname)s %(name)s: %(message)s",
- )
- await service.start()
- yield
- class OutboundCallRequest(BaseModel):
- sourceDn: str
- destination: str
- timeoutSec: int = 30
- app = FastAPI(
- title="3CX Telefonie Middleware",
- version="0.2.0",
- lifespan=lifespan,
- )
- app.add_middleware(
- CORSMiddleware,
- allow_origins=[
- x.strip()
- for x in settings.cors_origins.split(",")
- if x.strip()
- ],
- allow_credentials=False,
- allow_methods=["GET", "POST", "OPTIONS"],
- allow_headers=[
- "Authorization",
- "Content-Type",
- "Idempotency-Key",
- ],
- )
- @app.get("/phone/normalize")
- async def phone_normalize(
- number: str = Query(..., min_length=1),
- country: str = Query("DE", min_length=2, max_length=2),
- ):
- return normalize(number, country.upper())
- @app.get("/calls/active")
- async def active_calls():
- queue_dn = settings.threecx_queue_dn
- participants = service.get_active_participants()
- data = [
- normalize_call(participant, queue_dn)
- for participant in participants
- if participant.get("party_dn_type") == "Wexternalline"
- ]
- return {
- "data": data,
- "meta": {
- "totalItems": len(data)
- }
- }
- @app.post("/calls/outbound")
- async def outbound_call(
- request: OutboundCallRequest,
- authorization: str | None = Header(default=None),
- idempotency_key: str | None = Header(default=None, alias="Idempotency-Key"),
- ):
- if not authorized(authorization):
- return problem(
- 401,
- "Unauthorized",
- "Bearer-Token fehlt oder ist ungültig.",
- retryable=False,
- )
- if not idempotency_key or not idempotency_key.strip():
- return problem(
- 400,
- "Idempotency-Key fehlt",
- "Für ausgehende Anrufe ist ein Idempotency-Key erforderlich.",
- retryable=False,
- )
- idempotency_key = idempotency_key.strip()
- if len(idempotency_key) > 200:
- return problem(
- 400,
- "Idempotency-Key ungültig",
- "Der Idempotency-Key darf maximal 200 Zeichen lang sein.",
- retryable=False,
- )
- try:
- dn = request.sourceDn.strip()
- if not dn.isdigit():
- return problem(
- 422,
- "Ungültige Nebenstelle",
- "sourceDn muss eine numerische 3CX-Nebenstelle sein.",
- retryable=False,
- )
- if dn not in settings.monitored_dns:
- return problem(
- 403,
- "Nebenstelle nicht freigegeben",
- f"Die Nebenstelle {dn} ist für ausgehende API-Anrufe nicht freigegeben.",
- retryable=False,
- )
- raw_destination = request.destination.strip()
- # 1-4-stellige numerische Ziele sind interne 3CX-Nebenstellen.
- if raw_destination.isdigit() and 1 <= len(raw_destination) <= 4:
- normalized = {
- "raw": raw_destination,
- "e164": raw_destination,
- "valid": True,
- "type": "EXTENSION",
- }
- else:
- normalized = normalize(
- raw_destination,
- settings.phone_default_region,
- )
- if not normalized.get("valid"):
- return problem(
- 422,
- "Ungültige Rufnummer",
- "Die Zielrufnummer konnte nicht gültig normalisiert werden.",
- retryable=False,
- )
- request_fingerprint = hashlib.sha256(
- json.dumps(
- {
- "sourceDn": dn,
- "destination": normalized.get("e164"),
- "timeoutSec": request.timeoutSec,
- },
- sort_keys=True,
- separators=(",", ":"),
- ).encode("utf-8")
- ).hexdigest()
- existing = await service.repo.get_outbound_idempotency(idempotency_key)
- if existing:
- if existing["request_hash"] != request_fingerprint:
- return problem(
- 409,
- "Idempotency-Key bereits verwendet",
- "Der Idempotency-Key wurde bereits für eine andere Anfrage verwendet.",
- retryable=False,
- )
- if existing["status"] == "completed" and existing["response_json"]:
- return json.loads(existing["response_json"])
- return problem(
- 409,
- "Anruf wird bereits verarbeitet",
- "Für diesen Idempotency-Key läuft bereits ein Anrufvorgang.",
- retryable=True,
- )
- claimed = await service.repo.claim_outbound_idempotency(
- idempotency_key,
- request_fingerprint,
- datetime.now(timezone.utc).isoformat(),
- )
- if not claimed:
- existing = await service.repo.get_outbound_idempotency(idempotency_key)
- if existing and existing["request_hash"] != request_fingerprint:
- return problem(
- 409,
- "Idempotency-Key bereits verwendet",
- "Der Idempotency-Key wurde bereits für eine andere Anfrage verwendet.",
- retryable=False,
- )
- return problem(
- 409,
- "Anruf wird bereits verarbeitet",
- "Für diesen Idempotency-Key läuft bereits ein Anrufvorgang.",
- retryable=True,
- )
- dn_data = await service.cx.get_dn(dn)
- devices = dn_data.get("devices") or []
- if len(devices) == 0:
- return problem(
- 409,
- "Kein Gerät gefunden",
- f"Für Nebenstelle {dn} wurde kein 3CX-Gerät gefunden.",
- retryable=True,
- )
- hint = settings.outbound_device_hints.get(dn)
- if not hint:
- return problem(
- 409,
- "Kein Hardware-Gerät konfiguriert",
- f"Für Nebenstelle {dn} ist kein Hardware-Telefon für ausgehende API-Anrufe konfiguriert.",
- retryable=False,
- )
- matching_devices = [
- device
- for device in devices
- if hint.lower() in str(device.get("user_agent", "")).lower()
- ]
- if len(matching_devices) == 0:
- return problem(
- 409,
- "Hardware-Gerät nicht gefunden",
- f"Für Nebenstelle {dn} wurde kein 3CX-Gerät passend zu '{hint}' gefunden.",
- retryable=True,
- )
- if len(matching_devices) > 1:
- return problem(
- 409,
- "Hardware-Gerät nicht eindeutig",
- f"Für Nebenstelle {dn} wurden mehrere Geräte passend zu '{hint}' gefunden.",
- retryable=False,
- )
- device = matching_devices[0]
- device_id = (
- device.get("device_id")
- or device.get("deviceId")
- or device.get("id")
- )
- if not device_id:
- return problem(
- 502,
- "Ungültige 3CX-Geräteantwort",
- "3CX hat kein verwendbares Gerätekennzeichen geliefert.",
- retryable=True,
- )
- result = await service.cx.make_call(
- dn=dn,
- device_id=str(device_id),
- destination=normalized["e164"],
- timeout_sec=request.timeoutSec,
- )
- response = {
- "status": "accepted",
- "sourceDn": dn,
- "deviceId": str(device_id),
- "destination": {
- "raw": normalized.get("raw"),
- "e164": normalized.get("e164"),
- "valid": normalized.get("valid"),
- "type": normalized.get("type"),
- },
- "threecx": result,
- }
- threecx_result = result.get("result") or result
- threecx_callid = threecx_result.get("callid")
- threecx_legid = threecx_result.get("legid")
- if threecx_callid is not None:
- started_at = datetime.now(timezone.utc).isoformat()
- await service.repo.insert_outbound_call(
- callid=threecx_callid,
- legid=threecx_legid,
- source_dn=dn,
- destination=normalized.get("e164"),
- phone=(
- normalized.get("e164")
- if normalized.get("type") != "EXTENSION"
- else None
- ),
- raw=normalized.get("raw"),
- started_at=started_at,
- )
- await service.repo.finish_outbound_idempotency(
- idempotency_key,
- json.dumps(response, ensure_ascii=False, separators=(",", ":")),
- )
- return response
- except httpx.HTTPStatusError as exc:
- await service.repo.release_outbound_idempotency(idempotency_key)
- status = exc.response.status_code
- if status in (401, 403):
- return problem(
- 502,
- "3CX-Authentifizierung fehlgeschlagen",
- "Die Middleware konnte den 3CX-Anruf nicht authentifizieren.",
- retryable=False,
- )
- if status == 404:
- return problem(
- 502,
- "3CX-Gerät oder Ziel nicht gefunden",
- "3CX hat die angeforderte Ressource nicht gefunden.",
- retryable=False,
- )
- return problem(
- 502,
- "3CX-Anruf fehlgeschlagen",
- f"3CX hat HTTP {status} zurückgegeben.",
- retryable=True,
- )
- except Exception:
- await service.repo.release_outbound_idempotency(idempotency_key)
- log.exception("Outbound call failed")
- return problem(
- 502,
- "Ausgehender Anruf fehlgeschlagen",
- "Der Anruf konnte über 3CX nicht ausgelöst werden.",
- retryable=True,
- )
- @app.get("/calls/cdr")
- async def calls_cdr(
- from_: str = Query(..., alias="from"),
- to: str = Query(...),
- limit: int = Query(500, ge=1, le=1000),
- offset: int = Query(0, ge=0),
- authorization: str | None = Header(default=None),
- ):
- if not authorized(authorization):
- return problem(
- 401,
- "Unauthorized",
- "Bearer-Token fehlt oder ist ungültig.",
- retryable=False,
- )
- try:
- data = await service.cx.get_call_log(
- period_from=from_,
- period_to=to,
- top=limit,
- skip=offset,
- )
- return {
- "data": data.get("value", []),
- "meta": {
- "totalItems": data.get("@odata.count"),
- "limit": limit,
- "offset": offset,
- },
- }
- except Exception:
- log.exception("3CX CDR query failed")
- return problem(
- 502,
- "3CX CDR-Abfrage fehlgeschlagen",
- "Die Call-History konnte von 3CX nicht abgerufen werden.",
- retryable=True,
- )
- @app.get("/calls/recording/{rec_id}")
- async def calls_recording(
- rec_id: int,
- authorization: str | None = Header(default=None),
- ):
- if not authorized(authorization):
- return problem(
- 401,
- "Unauthorized",
- "Bearer-Token fehlt oder ist ungültig.",
- retryable=False,
- )
- try:
- from fastapi.responses import Response
- content, content_type = await service.cx.download_recording(rec_id)
- return Response(
- content=content,
- media_type=content_type.split(";", 1)[0],
- headers={
- "Content-Disposition": (
- f'attachment; filename="recording-{rec_id}.wav"'
- )
- },
- )
- except Exception:
- log.exception(
- "3CX recording download failed: rec_id=%s",
- rec_id,
- )
- return problem(
- 502,
- "Recording konnte nicht geladen werden",
- "Das 3CX-Recording konnte nicht abgerufen werden.",
- retryable=True,
- )
- @app.get("/calls/{call_id}")
- async def call_detail(call_id: int):
- row = await service.repo.get_call(call_id)
- if row is None:
- return problem(
- 404,
- "Call nicht gefunden",
- f"Call {call_id} wurde nicht gefunden.",
- retryable=False,
- )
- return row
- @app.get("/calls")
- async def calls_history(
- limit: int = Query(100, ge=1, le=500),
- offset: int = Query(0, ge=0),
- status: str | None = Query(None),
- ):
- rows, total = await service.repo.list_calls(
- limit=limit,
- offset=offset,
- status=status,
- )
- return {
- "data": rows,
- "meta": {
- "totalItems": total,
- "limit": limit,
- "offset": offset,
- }
- }
- @app.get("/health")
- async def health():
- return {
- "status": "ok",
- "queue": settings.threecx_queue_dn,
- "version": "0.2.0",
- }
- @app.get("/missed-calls")
- async def missed_calls():
- return await service.repo.list_groups()
- @app.get("/missed-calls/groups")
- async def missed_call_groups():
- groups = await service.repo.list_groups(open_only=True)
- return {
- "data": groups,
- "meta": {
- "totalItems": len(groups),
- },
- }
- @app.get("/missed-calls/groups/{e164}")
- async def missed_call_group(e164: str):
- group = await service.repo.get_group(e164)
- if not group:
- return problem(
- 404,
- "Not found",
- "Die Rufnummer wurde nicht gefunden.",
- )
- return group
- @app.post("/missed-calls/groups/{e164}/resolve")
- async def resolve_group(
- e164: str,
- payload: ResolveRequest,
- authorization: str | None = Header(default=None),
- idempotency_key: str | None = Header(default=None),
- ):
- if not authorized(authorization):
- return problem(
- 401,
- "Unauthorized",
- "Bearer-Token fehlt oder ist ungültig.",
- )
- if not idempotency_key:
- return problem(
- 400,
- "Idempotency-Key required",
- "Idempotency-Key ist erforderlich.",
- field="Idempotency-Key",
- )
- previous = await service.repo.get_idempotency(idempotency_key, "resolve")
- if previous:
- return JSONResponse(
- content=json.loads(previous),
- status_code=200,
- )
- group = await service.repo.get_group(e164)
- if not group:
- return problem(
- 404,
- "Not found",
- "Die Rufnummer wurde nicht gefunden.",
- )
- result = await service.repo.resolve(
- e164,
- payload.until,
- payload.by,
- payload.note,
- )
- response_json = json.dumps(result)
- await service.repo.save_idempotency(
- idempotency_key,
- "resolve",
- response_json,
- )
- return result
- @app.post("/missed-calls/groups/{e164}/reopen")
- async def reopen_group(
- e164: str,
- payload: ReopenRequest,
- authorization: str | None = Header(default=None),
- idempotency_key: str | None = Header(default=None),
- ):
- if not authorized(authorization):
- return problem(
- 401,
- "Unauthorized",
- "Bearer-Token fehlt oder ist ungültig.",
- )
- if not idempotency_key:
- return problem(
- 400,
- "Idempotency-Key required",
- "Idempotency-Key ist erforderlich.",
- field="Idempotency-Key",
- )
- previous = await service.repo.get_idempotency(
- idempotency_key,
- "reopen",
- )
- if previous:
- return JSONResponse(
- content=json.loads(previous),
- status_code=200,
- )
- group = await service.repo.get_group(e164)
- if not group:
- return problem(
- 404,
- "Not found",
- "Die Rufnummer wurde nicht gefunden.",
- )
- result = await service.repo.reopen(
- e164,
- payload.by,
- )
- response_json = json.dumps(result)
- await service.repo.save_idempotency(
- idempotency_key,
- "reopen",
- response_json,
- )
- return result
|