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

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_resolution e video_in_preferred_framerate nos 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

  1. 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
  2. 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.