| import os |
| import sys |
| import json |
| import signal |
| import argparse |
| import uvicorn |
| from fastapi import FastAPI, WebSocket |
| from fastapi.responses import JSONResponse |
|
|
| app = FastAPI() |
|
|
| @app.websocket("/media") |
| async def media_endpoint(websocket: WebSocket): |
| await websocket.accept() |
| print("[*] WebSocket connection accepted") |
| |
| while True: |
| try: |
| |
| data = await websocket.receive_text() |
| message = json.loads(data) |
| |
| |
| print(f"\n[*] Received message: {json.dumps(message, indent=2)}") |
| |
| |
| required_fields = { |
| "timestamp": str, |
| "streamId": str, |
| "callerId": str, |
| "channelId": str, |
| "event": str, |
| "callDirection": str, |
| "did": str, |
| "callId": str, |
| "cid": str, |
| "extraParams": str |
| } |
| |
| |
| if (all(field in message and isinstance(message[field], field_type) |
| for field, field_type in required_fields.items()) and |
| message["event"] == "answer"): |
| |
| response = { |
| "event": "reverse-media", |
| "callId": message["callId"], |
| "streamId": message["streamId"], |
| "chunk": 1, |
| "chunk_durn_ms": 20, |
| "mediaFormat": { |
| "encoding": "PCM", |
| "sampleRate": 8000, |
| "channels": 1 |
| }, |
| "payload": "AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAP//AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA//8AAAAAAAAAAAAAAAAAAAAA//8AAAAAAAD//////////////////wAA/////wAAAAD/////AAD//wAAAAABAP//AAAAAAEAAAAAAAEAAQABAAAAAQACAAEAAAABAAMAAAD//wEAAQABAAEAAAABAAAAAgD+/wAA//8BAPr/7f/z/+r/7P/w/+H/1//a/+L/5f/j/+r/8/////b/5v/j//3/AgDt/9r/6f8DAAAA7//9/xYADQDs/+//BgAJAAUAAgAJAA4ADwARABgAIAAlADIAMwAxAEQAUQBZAFIASgA/ADgAOwBFAEUARwBRAFMAOwAvAEIATQBPAEMAQgBJAEQAOgA8AEEAPQA9ADUAHwAgACYALQAaABIAFQATAAQA9f8FAPP/7v/w/+v/8f/X/83/x//L/83/yf/Q/7n/v/+9/8P/wP++/8j/yf/B/67/uv/X/9T/0v/I/9n/3v/K/9D/0v/0//P/8f/O/9z/8f/4//r/7v8LAAQAAwANAAwADgAaABkACwAZACMAMAA3ACMAHAAnACkAKwAoADoAPQA5AD0ANQAuACMARgBEADAANgBAAD8AKwApADIAOwAnAAsAEAASABEAGwAWAA4A/f/6//3/AQD6/+7/9v/+//X/4f/k/+r/5//s//X/7P/a/9j/4P/v/+v/2v/k/+X/5P/i/9T/2v/r//3/+P/7/wAA+f8IAAIABgAKACIAFAD9/xIAGwAsACMANgArAC4AJQAYACEAMwBTAEMAIgApADQAMAAWAC0ARgBFACAAGAAnACUAHwAPACAAIwAVAAkABQDy//T/FQAVAO7/5//1//D/3f/d//j//P/a/8L/wv+9/8X/2f/h/87/s//A/7v/u//C/73/vv+8/7//vP+//8T/y//I/8D/vf/K/9N/1v/P/9b/2//T/9f/2v/f/+r/4v/X/9f/7f/6//D/7//q//P/8P/1//7/AgALAA0ACgAHAAkAEAAQAA4ADgAZABMABgAJABYAGAAGAAwAEwAXABAABwAHAAoAEAAVABIACAABAAIABgAGABAAEAAWAA0A/P/4/wQACAADAAcABwAAAPP/8f/+/woABgALAAAA6//q//3/BAD1/wEABAD3//b/AwAKAPz/+P8JABoAEAAOAB0AGAAOAAcADAAaACYALQAoABgAHAAoACgAJwAiABgALQAwAB8AEQAkADgAMQAbAAcAIgAmACEABQAVACcAGwAIAPT/BgAXABwA/P/x/+//9f/x//n////0/+r/3f/b/+D/6v/u/+n/3P/U/8b/yP/T/93/4f/g/9d/vv+7/8j/2//c/9b/z//C/7j/wP/W/9//5v/d/83/u//G/9N/3f/Z/97/3//T/9H/2//1/+X/3P/j//r/9v/k/+r////+/+n/8P8AAAQA9f/1//v/BAD///v/AwD7/+z/9f8DAAAA+v/3//H/6v/t//P/9//y//f/9P/k/9j/4f/2/+7/6f/f/+L/6f/w/+f/3v/m/+P/5f/g/+f/7f/1/+X/4v/s//z/+//x//r/BgAIAPz/AAACAAoAHgAhABEADAAdACEAGwAdACkAMwAjABQAFwAmADQAOgBBADgAMgAwADQAOwBAAEQARABAADkAOwBDAEYAQABDAEIAPQAzADIANQA5ADMAMAAsACQAIAAhABoAFAAXAB0AEgAIAAkADQANAAcACQAOABAA+f/p/+f/+/8EAAAA/f/0//D/4f/m//T/9P/j/9z/3//Z/9b/3//u/w==" |
| } |
| await websocket.send_json(response) |
| print(f"[*] Sent response for answer event from callId: {message['callId']}") |
| |
| except Exception as e: |
| print(f"[!] Error: {str(e)}") |
| break |
|
|
| |
| @app.get("/health") |
| async def health_check_get(): |
| return JSONResponse(content={"message": "GET OK"}, status_code=200) |
|
|
| @app.post("/health") |
| async def health_check_post(): |
| return JSONResponse(content={"message": "POST OK"}, status_code=200) |
|
|
| def signal_handler(sig, frame): |
| print("\n[*] Shutting down server...") |
| sys.exit(0) |
|
|
| if __name__ == "__main__": |
| parser = argparse.ArgumentParser(description='FastAPI WebSocket Server') |
| parser.add_argument('--port', type=int, default=8000, help='Port to run the server on') |
| args = parser.parse_args() |
|
|
| signal.signal(signal.SIGINT, signal_handler) |
|
|
| print(f"[*] Server running at ws://localhost:{args.port}/media") |
| uvicorn.run("app:app", host="0.0.0.0", port=args.port, log_level="info") |