Skip to content

Commit a22022b

Browse files
authored
Merge pull request #12 from pipecat-ai/aleix/flush-before-close
WhiskerObserver: flush data before closing and finish task
2 parents ec24427 + a33376e commit a22022b

2 files changed

Lines changed: 38 additions & 15 deletions

File tree

CHANGELOG.md

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
2121

2222
### Fixed
2323

24+
- Fixed an issue where the buffer data was not being sent to the client when
25+
closing.
26+
27+
- Fixed an issue that was causing data to be written to the output file after
28+
the file was already closed.
29+
2430
- Fixed an issue that could cause `WhiskerObserver` to crash if the given
2531
pipeline didn't have a previous processor.
2632

pipecat/src/pipecat_whisker/observer.py

Lines changed: 32 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -146,7 +146,7 @@ def __init__(
146146
self._id = 0
147147
self._client = None
148148
self._server_future = asyncio.get_running_loop().create_future()
149-
self._server_task = asyncio.create_task(self._start_task_handler())
149+
self._server_task = asyncio.create_task(self._server_task_handler())
150150
self._send_task = asyncio.create_task(self._send_task_handler())
151151
self._send_queue = asyncio.Queue()
152152
self._batch = []
@@ -158,16 +158,13 @@ async def cleanup(self):
158158
"""Clean up resources and close the Whisker server."""
159159
await super().cleanup()
160160

161-
await self._maybe_close_file()
161+
await self._stop_send_task()
162162

163-
if self._client:
164-
await self._client.close(reason="Whisker shutting down")
163+
await self._close_client()
165164

166-
if not self._server_future.done():
167-
self._server_future.set_result(None)
165+
await self._stop_server()
168166

169-
if self._server_task:
170-
await self._server_task
167+
await self._maybe_close_file()
171168

172169
async def on_process_frame(self, data: FrameProcessed):
173170
"""Handle frame processing events.
@@ -187,7 +184,7 @@ async def on_push_frame(self, data: FramePushed):
187184
if not isinstance(data.frame, self._exclude_frames):
188185
await self._send_push_frame(data)
189186

190-
async def _start_task_handler(self):
187+
async def _server_task_handler(self):
191188
"""Start the Whisker server and handle incoming connections.
192189
193190
This method runs in a separate task and manages the websocket server lifecycle.
@@ -198,11 +195,18 @@ async def _start_task_handler(self):
198195
# Queue initial pipeline structure
199196
await self._send_pipeline()
200197

201-
async with serve(self._server_handler, self._host, self._port):
198+
async with serve(self._client_handler, self._host, self._port):
202199
logger.debug(f"ᓚᘏᗢ Whisker running at ws://{self._host}:{self._port}")
203200
await self._server_future
204201

205-
async def _server_handler(self, client):
202+
async def _stop_server(self):
203+
if not self._server_future.done():
204+
self._server_future.set_result(None)
205+
206+
if self._server_task:
207+
await self._server_task
208+
209+
async def _client_handler(self, client):
206210
"""Handle a new Whisker client connection.
207211
208212
Args:
@@ -227,6 +231,10 @@ async def _server_handler(self, client):
227231
logger.debug("ᓚᘏᗢ Whisker: client disconnected")
228232
await self._reset_client()
229233

234+
async def _close_client(self):
235+
if self._client:
236+
await self._client.close(reason="Whisker shutting down")
237+
230238
async def _reset_client(self):
231239
self._client = None
232240

@@ -243,18 +251,27 @@ async def _maybe_close_file(self):
243251

244252
async def _send_task_handler(self):
245253
"""Handle sending batched messages to the client."""
246-
while True:
254+
running = True
255+
while running:
247256
try:
248257
data, flush = await asyncio.wait_for(self._send_queue.get(), timeout=0.5)
249258

250-
self._batch.append(data)
259+
if data:
260+
self._batch.append(data)
251261

252262
await self._maybe_send_batch(flush=flush)
253263

254264
self._send_queue.task_done()
265+
266+
running = data is not None
255267
except asyncio.TimeoutError:
256268
await self._maybe_send_batch(flush=True)
257269

270+
async def _stop_send_task(self):
271+
await self._queue_data(None, True)
272+
if self._send_task:
273+
await self._send_task
274+
258275
async def _maybe_send_batch(self, *, flush: bool = False):
259276
"""Send batched messages to the client.
260277
@@ -418,9 +435,9 @@ async def _send_push_frame(self, data: FramePushed):
418435

419436
await self._queue_data(msg_packed)
420437

421-
async def _queue_data(self, msg: bytes, flush: bool = False):
438+
async def _queue_data(self, msg: Optional[bytes], flush: bool = False):
422439
await self._send_queue.put((msg, flush))
423-
if self._file:
440+
if self._file and msg:
424441
await self._file.write(msg)
425442

426443
async def _send(self, msg: bytes):

0 commit comments

Comments
 (0)