Vonage Video Transport para Pipecat
O Vonage Video Transport para Pipecat permite que você desenvolva aplicativos baseados em IA que participam de forma integrada das sessões da Video API da Vonage. Esse protocolo de transporte permite que você receba áudio e vídeo dos participantes da sessão e envie o áudio e o vídeo processados de volta para a sessão em tempo real.
O Pipecat é uma estrutura para a criação de aplicativos de IA conversacional de voz e multimodal. O Vonage Video Transport para Pipecat conecta o pipeline de processamento de mídia do Pipecat às sessões da Video API da Vonage, possibilitando uma ampla gama de casos de uso:
- Assistentes de IA de voz e vídeo em tempo real
- Serviços de transcrição e tradução em tempo real
- Gravação e análise de chamadas
- Processamento de efeitos de áudio e vídeo
- Moderação automatizada e filtragem de conteúdo
- Processamento e manipulação personalizados de mídia
O transport cuida da conversão de formatos de áudio e vídeo, do gerenciamento de sessões e da conectividade WebRTC, permitindo que você se concentre na criação da lógica do seu aplicativo.
Esta página inclui as seguintes seções:
- Introdução
- Requisitos
- Configuração básica
- Configuração do transporte
- Processamento de áudio e vídeo
- Assinatura de streaming
- Assinantes individuais de serviços de áudio
- Legendas
- Gerenciamento de sessões
- Integração de pipelines
- Melhores práticas
- Problemas conhecidos
Introdução
O Vonage Video Transport para Pipecat permite que você desenvolva aplicativos baseados em IA que participam de forma integrada das sessões da Video API da Vonage, oferecendo suporte a casos de uso tanto de áudio quanto de vídeo. O código está disponível em https://github.com/Vonage/pipecat.
Requisitos
Para utilizar esse transporte, você precisará da biblioteca Python Vonage Video Connector, que funciona nas plataformas Linux AMD64 e ARM64.
Configuração básica
Parâmetros de autenticação e sessão
Para usar o Vonage Video Transport para o Pipecat, você precisa de:
- ID do aplicativo - Seu identificador de aplicativo da Video API da Vonage
- ID da sessão - O ID da sessão da Video API à qual você deseja se juntar
- Token - Um token de participante válido para a sessão
Esses parâmetros podem ser obtidos no painel da Video API da Vonage ou gerados usando os SDKs de servidor da Video API da Vonage.
Inicializar o transporte
from pipecat.transports.vonage.video_connector import (
VonageVideoConnectorTransport,
VonageVideoConnectorTransportParams,
)
# Initialize the transport
transport = VonageVideoConnectorTransport(
application_id="your_application_id",
session_id="your_session_id",
token="your_participant_token",
params=VonageVideoConnectorTransportParams()
)
Configuração de transporte
Parâmetros básicos de áudio e vídeo
Configure o transporte de acordo com suas necessidades específicas de áudio e vídeo:
from pipecat.audio.vad.silero import SileroVADAnalyzer
transport_params = VonageVideoConnectorTransportParams(
audio_in_enabled=True, # Enable receiving audio
audio_out_enabled=True, # Enable sending audio
video_in_enabled=True, # Enable receiving video
video_out_enabled=True, # Enable sending video
publisher_name="AI Assistant", # Name for the published stream
audio_in_sample_rate=16000, # Input sample rate (Hz)
audio_in_channels=1, # Input channels
audio_out_sample_rate=24000, # Output sample rate (Hz)
audio_out_channels=1, # Output channels
video_out_width=1280, # Output video width
video_out_height=720, # Output video height
video_out_framerate=30, # Output video framerate
video_out_color_format="RGB", # Output video color format
vad_analyzer=SileroVADAnalyzer(), # Voice activity detection
audio_in_auto_subscribe=True, # Auto-subscribe to audio streams
video_in_auto_subscribe=False, # Auto-subscribe to video streams
captions_in_enabled=False, # Enable receiving captions
captions_in_auto_subscribe=False, # Auto-subscribe to caption streams
video_in_preferred_resolution=(640, 480), # Preferred input resolution
video_in_preferred_framerate=15, # Preferred input framerate
publisher_enable_opus_dtx=False, # Enable Opus DTX
session_enable_migration=False, # Enable session migration
video_connector_log_level="INFO", # Log level
clear_buffers_on_interruption=True, # Clear buffers on interruption
)
transport = VonageVideoConnectorTransport(
application_id,
session_id,
token,
params=transport_params
)
Detecção de Atividade de Voz (VAD)
É recomendável utilizar esse transporte com a Detecção de Atividade de Voz para otimizar o processamento de áudio:
from pipecat.audio.vad.silero import SileroVADAnalyzer
# Configure VAD for better audio processing
vad = SileroVADAnalyzer()
transport_params = VonageVideoConnectorTransportParams(
vad_analyzer=vad,
# ... other parameters
)
O VAD ajuda a reduzir o processamento desnecessário, detectando quando há fala no fluxo de áudio.
Limpeza do buffer em caso de interrupções
O clear_buffers_on_interruption O parâmetro determina se os buffers de mídia são
esvaziados automaticamente quando um quadro de interrupção é recebido no pipeline.
transport_params = VonageVideoConnectorTransportParams(
clear_buffers_on_interruption=True, # Default: True
# ... other parameters
)
Quando ativar (True, padrão):
- Applications de IA conversacional nos quais se deseja interromper a reprodução imediatamente quando o usuário interrompe
- Assistentes de voz interativos que precisam responder rapidamente às solicitações dos usuários
- Applications nos quais dados de áudio/vídeo desatualizados devem ser descartados para manter a interação em tempo real
- Situações em que minimizar a latência é mais importante do que concluir a reprodução de mídia
Quando desativar (False):
- Applications de gravação ou streaming nos quais você deseja preservar todo o conteúdo multimídia
- Applications que precisam concluir a reprodução de informações importantes, mesmo que sejam interrompidos
- Cenários de processamento em lote nos quais os meios devem ser processados sequencialmente, sem interrupção
- Casos de uso em que você está implementando uma lógica personalizada de tratamento de interrupções
Processamento de áudio e vídeo
Processamento de entradas de áudio e vídeo
O transporte converte automaticamente o áudio e o vídeo recebidos da sessão de vídeo da Vonage para os formatos de mídia internos do Pipecat:
# Audio and video input is handled automatically by the transport
# Incoming audio frames are converted to AudioRawFrame format
# Incoming video frames are converted to ImageRawFrame format
pipeline = Pipeline([
transport.input(), # Receives audio and video from Vonage session
# ... your AI processing pipeline
])
Geração de saídas de áudio e vídeo
Enviar áudio e vídeo de volta para a sessão do Vonage Video:
# Audio and video output is sent automatically through the pipeline
pipeline = Pipeline([
# ... your AI processing pipeline
transport.output(), # Sends audio and video to Vonage session
])
Assinatura de streaming
Quando o transport assina os fluxos dos participantes da sessão, ele gera quadros do Pipecat que seu pipeline pode processar. O comportamento difere entre áudio e vídeo:
Transmissões de vídeo: O transporte gera quadros de vídeo individuais para cada fluxo assinado, identificados pelo ID do fluxo. Isso permite que você processe vídeos de diferentes participantes separadamente em seu pipeline.
Fluxos de áudio: Por padrão, o transporte gera quadros de áudio com todos os fluxos de áudio assinados misturados. O áudio de todos os participantes é combinado em um único fluxo de áudio que seu pipeline recebe. Para casos de uso que exijam áudio por participante, consulte Assinantes individuais de serviços de áudio.
Por padrão, o transporte assina automaticamente os fluxos com base no audio_in_auto_subscribe
e video_in_auto_subscribe parâmetros. Você também pode controlar manualmente quais fluxos deseja assinar
para obter um controle mais detalhado.
Inscrição manual no feed
Se você precisar de mais controle sobre quais feeds deseja assinar, pode desativar a assinatura automática e assinar manualmente feeds específicos:
from pipecat.transports.vonage.video_connector import SubscribeSettings
# Disable auto-subscription
transport_params = VonageVideoConnectorTransportParams(
audio_in_auto_subscribe=False,
video_in_auto_subscribe=False,
# ... other parameters
)
# Manually subscribe when a participant joins
@transport.event_handler("on_participant_joined")
async def on_participant_joined(transport, data):
stream_id = data['streamId']
logger.info(f"Participant joined with stream {stream_id}, subscribing...")
await transport.subscribe_to_stream(
stream_id,
SubscribeSettings(
subscribe_to_audio=True,
subscribe_to_video=True,
preferred_resolution=(640, 480),
preferred_framerate=15
)
)
Quando usar a assinatura automática (audio_in_auto_subscribe=True, video_in_auto_subscribe=True, padrão):
- Applications simples nos quais você deseja receber todos os fluxos de todos os participantes
- Assistentes de voz ou de vídeo que precisam interagir com todos os participantes da sessão
- Applications de gravação ou monitoramento que devem capturar todos os participantes
- Casos de uso em que minimizar a complexidade do código é mais importante do que a assinatura seletiva
- Applications where all participants should be treated equally
Quando usar a assinatura manual (audio_in_auto_subscribe=False, video_in_auto_subscribe=False):
- Aplicativos que precisam realizar assinaturas seletivas com base nos metadados dos participantes ou na lógica da sessão
- Situações em que você deseja otimizar a largura de banda assinando apenas transmissões específicas
- Casos de uso que exigem configurações personalizadas de assinatura para cada participante (níveis de qualidade diferentes)
- Applications que precisam validar ou autenticar os participantes antes da inscrição
- Cenários complexos envolvendo várias partes, nos quais se deseja um controle detalhado sobre quais fluxos receber
Controle da qualidade do vídeo com transmissão simultânea
Ao assinar transmissões de vídeo, você pode controlar a qualidade do vídeo que recebe usando o
preferred_resolution e preferred_framerate parâmetros. Esses parâmetros são particularmente
úteis quando o editor está enviando transmissões simultâneas (várias camadas de qualidade).
Transmissões simultâneas contêm várias camadas espaciais e temporais:
- Camadas espaciais: Diferentes resoluções (por exemplo, 1280x720, 640x480, 320x240)
- Camadas temporais: Diferentes taxas de quadros (por exemplo, 30 fps, 15 fps, 7,5 fps)
Ao especificar a resolução e a taxa de quadros de sua preferência, você pode otimizar o uso da largura de banda e os requisitos de processamento:
# Subscribe to a lower quality stream for bandwidth efficiency
await transport.subscribe_to_stream(
stream_id,
SubscribeSettings(
subscribe_to_video=True,
preferred_resolution=(320, 240), # Request low spatial layer
preferred_framerate=15 # Request lower temporal layer
)
)
# Subscribe to a high quality stream for better visual fidelity
await transport.subscribe_to_stream(
stream_id,
SubscribeSettings(
subscribe_to_video=True,
preferred_resolution=(1280, 720), # Request high spatial layer
preferred_framerate=30 # Request higher temporal layer
)
)
Observações importantes:
- Se a emissora não oferecer transmissão simultânea ou se a versão solicitada não estiver disponível, o servidor fornecerá a qualidade disponível mais próxima
- Resoluções e taxas de quadros mais baixas reduzem o consumo de largura de banda e a sobrecarga de processamento
- Essas configurações afetam apenas a assinatura de vídeo; elas não determinam a qualidade de saída do editor
- Você também pode definir preferências globais usando
video_in_preferred_resolutionevideo_in_preferred_frameratenos parâmetros de transporte para transmissões com inscrição automática
Assinantes individuais de serviços de áudio
Beta: Atualmente, esse recurso está disponível na versão beta.
Além do fluxo de áudio misto padrão, o transporte oferece
UserAudioRawFrame quadros para cada participante inscrito. Cada quadro user_id O campo
é definido com o ID do stream do participante, permitindo que você distinga as fontes de áudio em seu
pipeline.
Tanto o caminho de áudio misto quanto o por assinante estão ativos simultaneamente — não é necessária nenhuma configuração adicional além de habilitar a entrada de áudio.
from pipecat.processors.frame_processor import FrameProcessor, FrameDirection
from pipecat.frames.frames import Frame, UserAudioRawFrame
class PerSubscriberAudioProcessor(FrameProcessor):
async def process_frame(self, frame: Frame, direction: FrameDirection):
await super().process_frame(frame, direction)
if isinstance(frame, UserAudioRawFrame):
logger.info(f"Audio from subscriber {frame.user_id}")
await self.push_frame(frame, direction)
Legendas
Beta: Atualmente, esse recurso está disponível na versão beta.
O protocolo de transporte permite receber legendas em tempo real dos participantes da sessão. As legendas são
transmitidas como TranscriptionFrame (final) ou InterimTranscriptionFrame (em andamento) e
enviado para o pipeline.
Nota: As legendas devem ser habilitadas no nível da sessão por meio da Video API da Vonage antes de poderem ser recebidas pelo protocolo de transporte. Consulte o Legendas em tempo real guia para mais detalhes.
Ativando as legendas
Para receber legendas, configure captions_in_enabled para True nos parâmetros de transporte. Você
também precisa se inscrever para receber legendas, seja automaticamente ou manualmente.
Assinatura de legendas automáticas
transport_params = VonageVideoConnectorTransportParams(
captions_in_enabled=True,
captions_in_auto_subscribe=True,
# ... other parameters
)
Assinatura de legendas manuais
Para um controle mais detalhado, desative a inscrição automática e inscreva-se nas legendas de cada transmissão:
from pipecat.transports.vonage.video_connector import SubscribeSettings
transport_params = VonageVideoConnectorTransportParams(
captions_in_enabled=True,
captions_in_auto_subscribe=False,
# ... other parameters
)
@transport.event_handler("on_participant_joined")
async def on_participant_joined(transport, data):
stream_id = data['streamId']
await transport.subscribe_to_stream(
stream_id,
SubscribeSettings(
subscribe_to_audio=True,
subscribe_to_captions=True,
)
)
Processamento de legendas no fluxo de trabalho
Os quadros de legenda incluem o texto transcrito e o ID de transmissão do participante como user_id, e
um carimbo de data e hora. TranscriptionFrame representa uma transcrição final, enquanto
InterimTranscriptionFrame representa uma transcrição em andamento que ainda pode sofrer alterações.
from pipecat.processors.frame_processor import FrameProcessor, FrameDirection
from pipecat.frames.frames import Frame, TranscriptionFrame, InterimTranscriptionFrame
class CaptionProcessor(FrameProcessor):
async def process_frame(self, frame: Frame, direction: FrameDirection):
await super().process_frame(frame, direction)
if isinstance(frame, TranscriptionFrame):
logger.info(f"[{frame.user_id}] Final: {frame.text}")
elif isinstance(frame, InterimTranscriptionFrame):
logger.info(f"[{frame.user_id}] Interim: {frame.text}")
await self.push_frame(frame, direction)
Gerenciamento de sessões
Eventos do ciclo de vida da sessão
Tratar eventos de entrada e saída da sessão:
# Handle when the transport joins the session
@transport.event_handler("on_joined")
async def on_joined(transport, data):
logger.info(f"Joined session {data['sessionId']}")
# Initialize your application state
await task.queue_frames([LLMMessagesFrame()])
# Handle when the transport leaves the session
@transport.event_handler("on_left")
async def on_left(transport):
logger.info("Left session")
# Clean up resources or save session data
# Handle session errors
@transport.event_handler("on_error")
async def on_error(transport, error):
logger.error(f"Session error: {error}")
Eventos com participantes
Acompanhar os participantes que entram e saem da sessão:
# Handle first participant joining (useful for triggering initial interactions)
@transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, data):
logger.info(f"First participant joined: stream {data['streamId']}")
# Start recording, initialize conversation, etc.
# Handle any participant joining
@transport.event_handler("on_participant_joined")
async def on_participant_joined(transport, data):
logger.info(f"Participant joined: stream {data['streamId']}")
# Update participant list, send greeting, etc.
# Handle participant leaving
@transport.event_handler("on_participant_left")
async def on_participant_left(transport, data):
logger.info(f"Participant left: stream {data['streamId']}")
# Update participant list, handle cleanup
Eventos de conexão do cliente
Monitorar as conexões individuais dos assinantes do stream:
# Handle when a subscriber successfully connects to a stream
@transport.event_handler("on_client_connected")
async def on_client_connected(transport, data):
logger.info(f"Client connected to stream {data['subscriberId']}")
# Handle when a subscriber disconnects from a stream
@transport.event_handler("on_client_disconnected")
async def on_client_disconnected(transport, data):
logger.info(f"Client disconnected from stream {data['subscriberId']}")
Integração de pipelines
Exemplo completo de pipeline
Veja a seguir como integrar o Vonage Video Transport para Pipecat a um pipeline completo de IA:
import asyncio
import os
from pipecat.audio.vad.silero import SileroVADAnalyzer
from pipecat.frames.frames import LLMRunFrame
from pipecat.pipeline.pipeline import Pipeline
from pipecat.pipeline.runner import PipelineRunner
from pipecat.pipeline.task import PipelineTask
from pipecat.processors.aggregators.openai_llm_context import OpenAILLMContext
from pipecat.services.aws_nova_sonic.aws import AWSNovaSonicLLMService
from pipecat.transports.vonage.video_connector import (
VonageVideoConnectorTransport,
VonageVideoConnectorTransportParams,
SubscribeSettings,
)
async def main():
# Configure transport
transport = VonageVideoConnectorTransport(
application_id=os.getenv("VONAGE_APPLICATION_ID"),
session_id=os.getenv("VONAGE_SESSION_ID"),
token=os.getenv("VONAGE_TOKEN"),
params=VonageVideoConnectorTransportParams(
audio_in_enabled=True,
audio_out_enabled=True,
publisher_name="AI Assistant",
audio_in_sample_rate=16000,
audio_out_sample_rate=24000,
vad_analyzer=SileroVADAnalyzer(),
audio_in_auto_subscribe=True,
)
)
# Set up AI service
llm = AWSNovaSonicLLMService(
secret_access_key=os.getenv("AWS_SECRET_ACCESS_KEY", ""),
access_key_id=os.getenv("AWS_ACCESS_KEY_ID", ""),
region=os.getenv("AWS_REGION", ""),
session_token=os.getenv("AWS_SESSION_TOKEN", ""),
)
# Create context for conversation
context = OpenAILLMContext(
messages=[
{"role": "system", "content": "You are a helpful AI assistant."},
]
)
context_aggregator = llm.create_context_aggregator(context)
# Build pipeline
pipeline = Pipeline([
transport.input(), # Audio from Vonage session
context_aggregator.user(), # Process user input
llm, # AI processing
transport.output(), # Audio back to session
])
task = PipelineTask(pipeline)
# Register event handlers
@transport.event_handler("on_joined")
async def on_joined(transport, data):
print(f"Joined session: {data['sessionId']}")
@transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, data):
print(f"First participant joined: {data['streamId']}")
@transport.event_handler("on_client_connected")
async def on_client_connected(transport, data):
await task.queue_frames([LLMRunFrame()])
await llm.trigger_assistant_response()
runner = PipelineRunner()
await runner.run(task)
if __name__ == "__main__":
asyncio.run(main())
Melhores práticas
Otimização de desempenho
-
Escolha taxas de amostragem adequadas:
- Use 16 kHz para reconhecimento de fala e a maioria dos serviços de IA
- Use 24 kHz ou mais para obter melhor qualidade na conversão de texto em fala
- Evite taxas de amostragem desnecessariamente altas, que aumentam a carga de processamento
-
Otimizar o processamento do pipeline:
- Mantenha os fluxos de processamento de IA eficientes para minimizar a latência
- Utilize tamanhos de quadro adequados e gerenciamento de buffer
- Considere utilizar o VAD para reduzir o processamento desnecessário
Depuração e monitoramento
Ativar registro em log:
from loguru import logger
# Pipecat uses loguru for logging
logger.enable("pipecat")
# Set Vonage Video Connector log level
transport_params = VonageVideoConnectorTransportParams(
video_connector_log_level="DEBUG", # DEBUG, INFO, WARNING, ERROR
# ... other parameters
)
Problemas conhecidos
- Em casos raros, a aplicação pode ficar momentaneamente sem resposta durante o processo de desligamento, principalmente quando um volume extremamente grande de transmissões está sendo simultaneamente retirado do ar, cancelado e desconectado à força na mesma sessão. Estamos trabalhando ativamente para resolver esse problema.