move SileroVAD processor to processors package
This commit is contained in:
@@ -17,6 +17,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
|||||||
|
|
||||||
### Changed
|
### Changed
|
||||||
|
|
||||||
|
- Moved `SileroVAD` audio processor to `processors.audio.vad`.
|
||||||
|
|
||||||
- Module `utils.audio` is now `audio.utils`. A new `resample_audio` function has
|
- Module `utils.audio` is now `audio.utils`. A new `resample_audio` function has
|
||||||
been added.
|
been added.
|
||||||
|
|
||||||
|
|||||||
@@ -9,7 +9,6 @@ import aiohttp
|
|||||||
import os
|
import os
|
||||||
import sys
|
import sys
|
||||||
|
|
||||||
from pipecat.audio.vad.silero import SileroVAD
|
|
||||||
from pipecat.frames.frames import LLMMessagesFrame
|
from pipecat.frames.frames import LLMMessagesFrame
|
||||||
from pipecat.pipeline.pipeline import Pipeline
|
from pipecat.pipeline.pipeline import Pipeline
|
||||||
from pipecat.pipeline.runner import PipelineRunner
|
from pipecat.pipeline.runner import PipelineRunner
|
||||||
@@ -18,6 +17,7 @@ from pipecat.processors.aggregators.llm_response import (
|
|||||||
LLMAssistantResponseAggregator,
|
LLMAssistantResponseAggregator,
|
||||||
LLMUserResponseAggregator,
|
LLMUserResponseAggregator,
|
||||||
)
|
)
|
||||||
|
from pipecat.processors.audio.vad.silero import SileroVAD
|
||||||
from pipecat.services.cartesia import CartesiaTTSService
|
from pipecat.services.cartesia import CartesiaTTSService
|
||||||
from pipecat.services.openai import OpenAILLMService
|
from pipecat.services.openai import OpenAILLMService
|
||||||
from pipecat.transports.services.daily import DailyParams, DailyTransport
|
from pipecat.transports.services.daily import DailyParams, DailyTransport
|
||||||
|
|||||||
@@ -8,16 +8,7 @@ import time
|
|||||||
|
|
||||||
import numpy as np
|
import numpy as np
|
||||||
|
|
||||||
from pipecat.audio.vad.vad_analyzer import VADAnalyzer, VADParams, VADState
|
from pipecat.audio.vad.vad_analyzer import VADAnalyzer, VADParams
|
||||||
from pipecat.frames.frames import (
|
|
||||||
AudioRawFrame,
|
|
||||||
Frame,
|
|
||||||
StartInterruptionFrame,
|
|
||||||
StopInterruptionFrame,
|
|
||||||
UserStartedSpeakingFrame,
|
|
||||||
UserStoppedSpeakingFrame,
|
|
||||||
)
|
|
||||||
from pipecat.processors.frame_processor import FrameDirection, FrameProcessor
|
|
||||||
|
|
||||||
from loguru import logger
|
from loguru import logger
|
||||||
|
|
||||||
@@ -171,75 +162,3 @@ class SileroVADAnalyzer(VADAnalyzer):
|
|||||||
# This comes from an empty audio array
|
# This comes from an empty audio array
|
||||||
logger.exception(f"Error analyzing audio with Silero VAD: {e}")
|
logger.exception(f"Error analyzing audio with Silero VAD: {e}")
|
||||||
return 0
|
return 0
|
||||||
|
|
||||||
|
|
||||||
class SileroVAD(FrameProcessor):
|
|
||||||
def __init__(
|
|
||||||
self,
|
|
||||||
*,
|
|
||||||
sample_rate: int = 16000,
|
|
||||||
vad_params: VADParams = VADParams(),
|
|
||||||
audio_passthrough: bool = False,
|
|
||||||
):
|
|
||||||
super().__init__()
|
|
||||||
|
|
||||||
self._vad_analyzer = SileroVADAnalyzer(sample_rate=sample_rate, params=vad_params)
|
|
||||||
self._audio_passthrough = audio_passthrough
|
|
||||||
|
|
||||||
self._processor_vad_state: VADState = VADState.QUIET
|
|
||||||
|
|
||||||
#
|
|
||||||
# FrameProcessor
|
|
||||||
#
|
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
if isinstance(frame, AudioRawFrame):
|
|
||||||
await self._analyze_audio(frame)
|
|
||||||
if self._audio_passthrough:
|
|
||||||
await self.push_frame(frame, direction)
|
|
||||||
else:
|
|
||||||
await self.push_frame(frame, direction)
|
|
||||||
|
|
||||||
#
|
|
||||||
# Handle interruptions
|
|
||||||
#
|
|
||||||
|
|
||||||
async def _handle_interruptions(self, frame: Frame):
|
|
||||||
if self.interruptions_allowed:
|
|
||||||
# Make sure we notify about interruptions quickly out-of-band.
|
|
||||||
if isinstance(frame, UserStartedSpeakingFrame):
|
|
||||||
logger.debug("User started speaking")
|
|
||||||
await self._start_interruption()
|
|
||||||
# Push an out-of-band frame (i.e. not using the ordered push
|
|
||||||
# frame task) to stop everything, specially at the output
|
|
||||||
# transport.
|
|
||||||
await self.push_frame(StartInterruptionFrame())
|
|
||||||
elif isinstance(frame, UserStoppedSpeakingFrame):
|
|
||||||
logger.debug("User stopped speaking")
|
|
||||||
await self._stop_interruption()
|
|
||||||
await self.push_frame(StopInterruptionFrame())
|
|
||||||
|
|
||||||
await self.push_frame(frame)
|
|
||||||
|
|
||||||
async def _analyze_audio(self, frame: AudioRawFrame):
|
|
||||||
# Check VAD and push event if necessary. We just care about changes
|
|
||||||
# from QUIET to SPEAKING and vice versa.
|
|
||||||
new_vad_state = self._vad_analyzer.analyze_audio(frame.audio)
|
|
||||||
if (
|
|
||||||
new_vad_state != self._processor_vad_state
|
|
||||||
and new_vad_state != VADState.STARTING
|
|
||||||
and new_vad_state != VADState.STOPPING
|
|
||||||
):
|
|
||||||
new_frame = None
|
|
||||||
|
|
||||||
if new_vad_state == VADState.SPEAKING:
|
|
||||||
new_frame = UserStartedSpeakingFrame()
|
|
||||||
elif new_vad_state == VADState.QUIET:
|
|
||||||
new_frame = UserStoppedSpeakingFrame()
|
|
||||||
|
|
||||||
if new_frame:
|
|
||||||
await self._handle_interruptions(new_frame)
|
|
||||||
|
|
||||||
self._processor_vad_state = new_vad_state
|
|
||||||
|
|||||||
0
src/pipecat/processors/audio/vad/__init__.py
Normal file
0
src/pipecat/processors/audio/vad/__init__.py
Normal file
91
src/pipecat/processors/audio/vad/silero.py
Normal file
91
src/pipecat/processors/audio/vad/silero.py
Normal file
@@ -0,0 +1,91 @@
|
|||||||
|
#
|
||||||
|
# Copyright (c) 2024, Daily
|
||||||
|
#
|
||||||
|
# SPDX-License-Identifier: BSD 2-Clause License
|
||||||
|
#
|
||||||
|
|
||||||
|
from pipecat.audio.vad.silero import SileroVADAnalyzer
|
||||||
|
from pipecat.audio.vad.vad_analyzer import VADParams, VADState
|
||||||
|
from pipecat.frames.frames import (
|
||||||
|
AudioRawFrame,
|
||||||
|
Frame,
|
||||||
|
StartInterruptionFrame,
|
||||||
|
StopInterruptionFrame,
|
||||||
|
UserStartedSpeakingFrame,
|
||||||
|
UserStoppedSpeakingFrame,
|
||||||
|
)
|
||||||
|
from pipecat.processors.frame_processor import FrameDirection, FrameProcessor
|
||||||
|
|
||||||
|
from loguru import logger
|
||||||
|
|
||||||
|
|
||||||
|
class SileroVAD(FrameProcessor):
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
sample_rate: int = 16000,
|
||||||
|
vad_params: VADParams = VADParams(),
|
||||||
|
audio_passthrough: bool = False,
|
||||||
|
):
|
||||||
|
super().__init__()
|
||||||
|
|
||||||
|
self._vad_analyzer = SileroVADAnalyzer(sample_rate=sample_rate, params=vad_params)
|
||||||
|
self._audio_passthrough = audio_passthrough
|
||||||
|
|
||||||
|
self._processor_vad_state: VADState = VADState.QUIET
|
||||||
|
|
||||||
|
#
|
||||||
|
# FrameProcessor
|
||||||
|
#
|
||||||
|
|
||||||
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
|
await super().process_frame(frame, direction)
|
||||||
|
|
||||||
|
if isinstance(frame, AudioRawFrame):
|
||||||
|
await self._analyze_audio(frame)
|
||||||
|
if self._audio_passthrough:
|
||||||
|
await self.push_frame(frame, direction)
|
||||||
|
else:
|
||||||
|
await self.push_frame(frame, direction)
|
||||||
|
|
||||||
|
#
|
||||||
|
# Handle interruptions
|
||||||
|
#
|
||||||
|
|
||||||
|
async def _handle_interruptions(self, frame: Frame):
|
||||||
|
if self.interruptions_allowed:
|
||||||
|
# Make sure we notify about interruptions quickly out-of-band.
|
||||||
|
if isinstance(frame, UserStartedSpeakingFrame):
|
||||||
|
logger.debug("User started speaking")
|
||||||
|
await self._start_interruption()
|
||||||
|
# Push an out-of-band frame (i.e. not using the ordered push
|
||||||
|
# frame task) to stop everything, specially at the output
|
||||||
|
# transport.
|
||||||
|
await self.push_frame(StartInterruptionFrame())
|
||||||
|
elif isinstance(frame, UserStoppedSpeakingFrame):
|
||||||
|
logger.debug("User stopped speaking")
|
||||||
|
await self._stop_interruption()
|
||||||
|
await self.push_frame(StopInterruptionFrame())
|
||||||
|
|
||||||
|
await self.push_frame(frame)
|
||||||
|
|
||||||
|
async def _analyze_audio(self, frame: AudioRawFrame):
|
||||||
|
# Check VAD and push event if necessary. We just care about changes
|
||||||
|
# from QUIET to SPEAKING and vice versa.
|
||||||
|
new_vad_state = self._vad_analyzer.analyze_audio(frame.audio)
|
||||||
|
if (
|
||||||
|
new_vad_state != self._processor_vad_state
|
||||||
|
and new_vad_state != VADState.STARTING
|
||||||
|
and new_vad_state != VADState.STOPPING
|
||||||
|
):
|
||||||
|
new_frame = None
|
||||||
|
|
||||||
|
if new_vad_state == VADState.SPEAKING:
|
||||||
|
new_frame = UserStartedSpeakingFrame()
|
||||||
|
elif new_vad_state == VADState.QUIET:
|
||||||
|
new_frame = UserStoppedSpeakingFrame()
|
||||||
|
|
||||||
|
if new_frame:
|
||||||
|
await self._handle_interruptions(new_frame)
|
||||||
|
|
||||||
|
self._processor_vad_state = new_vad_state
|
||||||
Reference in New Issue
Block a user