average-architect-28889
07/30/2025, 8:16 AMtall-airline-16756
07/30/2025, 8:16 AMtall-airline-16756
07/30/2025, 8:17 AMBlockOnClose parameter to ensure the sink processes all messages in the write queue before closing. This prevents dropped messages and ensures other participants are notified of departures.
• Agent Sessions: Use the drain() method before closing. This ensures all activities are completed, including uninterruptible speech, and detaches inputs and outputs properly.
• Workers: Call the drain() method with an optional timeout. This sets the worker status to full, waits for all processes to complete, and applies a timeout to prevent indefinite blocking.
• LiveKit Server: During server shutdown, LiveKit enters a draining mode where active rooms continue to run, new participants are rejected for new rooms, and the server waits for all participants to disconnect before shutting down.
Best Practices:
• Always call drain() before close().
• Set appropriate timeouts to avoid indefinite blocking.
• Check if the component is already in a draining state to avoid duplicate operations.
• Enable BlockOnClose for signal connections to ensure message delivery.
• Ensure workers gracefully drain active sessions before shutting down.
These steps will help ensure tasks are properly drained before closing, preventing data loss and ensuring smooth operations.
Sources: pkg/routing/messagechannel_test.go | livekit-agents/livekit/agents/ipc/job_proc_lazy_main.py | agents/src/worker.ts | Distributed Setup | LiveKit Docs | Worker lifecycle | LiveKit Docsaverage-architect-28889
07/30/2025, 11:52 AMtall-airline-16756
07/30/2025, 11:53 AMasync drain(timeout?: number) {
if (this.#draining) {
return;
}
this.#<http://logger.info|logger.info>('draining worker');
this.#draining = true;
this.event.emit(
'worker_msg',
new WorkerMessage({
message: {
case: 'updateWorker',
value: {
status: WorkerStatus.WS_FULL,
},
},
}),
);
const joinJobs = async () => {
return Promise.all(
this.#procPool.processes.map((proc) => {
if (!proc.runningJob) {
proc.close();
}
return proc.join();
}),
);
};
let timer: NodeJS.Timeout | undefined;
if (timeout) {
timer = setTimeout(() => {
throw new WorkerError('timed out draining');
}, timeout);
}
await joinJobs().then(() => {
if (timeout) {
clearTimeout(timer);
}
});
}
2. Python Agent Session:
async def drain(self) -> None:
if self._activity is None:
raise RuntimeError("AgentSession isn't running")
await self._activity.drain()
3. Go Egress AppWriter:
func (w *AppWriter) Drain(force bool) {
w.draining.Once(func() {
w.logger.Debugw("draining")
if force || !w.active.Load() {
w.endStream.Break()
} else {
time.AfterFunc(w.conf.Latency.PipelineLatency, func() { w.endStream.Break() })
}
})
<-w.finished.Watch()
}
4. Python Room Tasks:
async def _drain_rpc_invocation_tasks(self) -> None:
if self._rpc_invocation_tasks:
for task in self._rpc_invocation_tasks:
task.cancel()
await asyncio.gather(*self._rpc_invocation_tasks, return_exceptions=True)
async def _drain_data_stream_tasks(self) -> None:
if self._data_stream_tasks:
for task in self._data_stream_tasks:
task.cancel()
await asyncio.gather(*self._data_stream_tasks, return_exceptions=True)
These examples are like a friendly roadmap for implementing the drain method across different components. Think of it as a gentle way to wrap up tasks and shut things down gracefully. Hope this helps you out! 👍
Sources: agents/src/cli.ts | livekit-agents/livekit/agents/voice/agent_session.py | pkg/pipeline/source/sdk/appwriter.go | livekit-rtc/livekit/rtc/room.py