Repository navigation
Expand file tree
/
Copy pathserver.py
More file actions
150 lines (117 loc) · 4.7 KB
/
Copy pathserver.py
File metadata and controls
150 lines (117 loc) · 4.7 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
import os
import sys
from datetime import UTC, datetime
from fastapi import FastAPI, WebSocket
from fastapi.middleware.cors import CORSMiddleware
from intake_bot.bot import run_bot
from intake_bot.nodes.nodes import node_initial
from intake_bot.utils.ev import ev_is_true, get_ev
from loguru import logger
from pipecat.audio.mixers.base_audio_mixer import BaseAudioMixer
from pipecat.serializers.protobuf import ProtobufFrameSerializer
from pipecat.transports.websocket.fastapi import (
FastAPIWebsocketParams,
FastAPIWebsocketTransport,
)
logger.remove(0)
# Suppress noisy pipecat DEBUG logs from turn-detection internals.
_NOISY_PIPECAT_MODULES = {
"pipecat.turns.user_start",
"pipecat.audio.turn.smart_turn",
}
def _log_filter(record):
if record["level"].name == "DEBUG":
name = record["name"] or ""
for prefix in _NOISY_PIPECAT_MODULES:
if name.startswith(prefix):
return False
return True
logger.add(sys.stderr, level=get_ev("LOG_LEVEL", "INFO"), filter=_log_filter)
if ev_is_true("LOG_TO_FILE"):
os.makedirs("logs", exist_ok=True)
logger.add("logs/server.log", level=get_ev("LOG_LEVEL", "INFO"), filter=_log_filter)
def generate_call_id() -> str:
return datetime.now(UTC).strftime("%Y%m%dT%H%M%S%fZ")
class SilenceMixer(BaseAudioMixer):
"""Passthrough audio mixer that maintains a continuous WebSocket audio stream.
The FastAPIWebsocketTransport calls ``mix()`` on every audio output cycle,
including cycles where the bot is silent (listening or processing). Without
a mixer the transport only emits frames when TTS audio is present, so the
client receives nothing during silence. Many WebSocket audio clients
(browsers, the Python test client) expect a steady byte stream and will
stall, mis-time playback, or drop the connection if the stream goes quiet.
This mixer solves the problem with the simplest possible implementation:
pass every audio buffer through unchanged. No actual mixing is required
because only one audio source (TTS) is in play; the mixer is registered
purely to opt in to the continuous-output behaviour of the transport.
"""
async def start(self, sample_rate: int):
pass
async def stop(self):
pass
async def process_frame(self, frame):
pass
async def mix(self, audio: bytes) -> bytes:
return audio
def _get_user_idle_timeout_secs(websocket: WebSocket, call_id: str) -> float | None:
raw_timeout = websocket.query_params.get("idle_timeout_secs")
if raw_timeout is None:
raw_timeout = get_ev("WEBSOCKET_USER_IDLE_TIMEOUT_SECS", "").strip()
if raw_timeout is None or raw_timeout == "":
if call_id.startswith("ws-test"):
raw_timeout = get_ev("WEBSOCKET_TEST_USER_IDLE_TIMEOUT_SECS", "45.0")
if not raw_timeout:
return None
try:
timeout_secs = float(raw_timeout)
except ValueError:
logger.warning(
f"""Ignoring invalid websocket idle timeout value: {raw_timeout!r}"""
)
return None
if timeout_secs <= 0:
logger.warning(
f"""Ignoring non-positive websocket idle timeout value: {raw_timeout!r}"""
)
return None
return timeout_secs
app = FastAPI()
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=True,
allow_methods=["*"],
allow_headers=["*"],
)
@app.get("/healthz")
async def healthz():
return {"status": "ok"}
@app.websocket("/ws")
async def websocket_endpoint(websocket: WebSocket):
await websocket.accept()
caller_phone_number = websocket.query_params.get("caller_phone_number", "")
call_id = websocket.query_params.get("call_id") or generate_call_id()
user_idle_timeout_secs = _get_user_idle_timeout_secs(websocket, call_id)
transport = FastAPIWebsocketTransport(
websocket=websocket,
params=FastAPIWebsocketParams(
audio_in_enabled=True,
audio_out_enabled=True,
add_wav_header=False,
serializer=ProtobufFrameSerializer(),
audio_out_mixer=SilenceMixer(),
),
)
async def configure_websocket_transport(transport, task, flow_manager, call_id):
@transport.event_handler("on_client_connected")
async def on_client_connected(transport, client):
logger.info(f"""WebSocket client connected for call {call_id}""")
await flow_manager.initialize(node_initial())
await run_bot(
transport=transport,
call_id=call_id,
caller_phone_number=caller_phone_number,
handle_sigint=False,
configure_transport=configure_websocket_transport,
user_idle_timeout_secs=user_idle_timeout_secs,
)