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