Pular para conteúdo

Orchestrator

Papel e localização

O Orchestrator converte uma proposta do agente em alterações válidas da conversa e coordena a entrega da resposta. Há duas partes: a Lambda src/orchestrator/app.py administra transporte, retries e atualidade da mensagem; src/core/application/process_message_use_case.py administra sessão, agentes e ações. O caso de uso não envia a resposta final ao WhatsApp: retorna ProcessOutcome para a Lambda.

Símbolo Contrato e consumidor
lambda_handler(event, context) Recebe Records SQS; retorna {"batchItemFailures": [...]}
_process_record(...) Extrai telefone, drena receipts, tenta entrega pendente ou buffer legado
_process_pending_turns(...) Processa receipts em ordem DynamoDB e constrói os guards
process_message(deps, *, phone_number, text, debounced_message_ids, inbound_message_id=None, inbound_received_at=None, whatsapp_instance=..., turn_guard=None) Resolve sessão e executa um turno; retorna ProcessOutcome
ProcessOutcome reply: str | None, state_after: SessionState, superseded: bool=False
TurnGuard Callbacks is_current, try_commit, try_persist, discard_stale_without_session
_agent_loop(...) Seleciona agente, invoca, faz parse, dispatch e aplica ActionKind

Dependencies, de src/core/integrations/factory.py, fornece repositories, clientes, KB, settings e invocador. O AgentInvoker é um Protocol; testes conseguem substituir toda a geração sem modificar o caso de uso.

Resolução da sessão e histórico

_resolve_active_session() consulta o pointer ACTIVE por telefone. Na ausência de sessão cria Session.new(); se a sessão atual é CLOSED ou OUT_OF_SCOPE, cria uma nova vinculada por previous_session_id. _create_session() pode carregar o nome cadastrado anteriormente em recognized_name. Não transfere o cadastro inteiro nem todo o histórico.

O inbound é registrado com seu ID original. HUMAN retorna sem IA; HANDOFF_PENDING registra e devolve confirmação fixa de recebimento. Nos estados do bot, zera context_enrichment_count para o turno e converte NEW → TRIAGE. _session_history() chama history_for_session(..., limit=MAX_HISTORY_MESSAGES), evitando que outro assunto ou resposta humana antiga reapareça no prompt.

Seleção determinística: TRIAGE → triagem, QUALIFYING → qualificacao, DATA_COLLECT/CONFIRM → coleta. Qualquer outro estado na seleção lança ApplicationError.

Aplicação das ações

ActionKind Aplicação concreta
REPLY_AND_STAY Merge de campos, atualização de classificação, transição se necessária, persistência e log outbound
NEXT_AGENT_SAME_TURN Merge e troca de estado/agente; chama o próximo sem repetir entrada do usuário
CONTEXT_ENRICH Guarda ID do package, incrementa contador e chama novamente o agente
SUBJECT_CHANGE_CONFIRMATION Guarda original_text em pending_subject_change_text e responde pedindo confirmação
SUBJECT_CHANGE Resumo, fechamento, nova sessão ligada e replay da solicitação original
HANDOFF Guarda motivo/contexto, resumo, HANDOFF_PENDING, persistência e notificação
OUT_OF_PORTFOLIO Cria demanda, fecha em OUT_OF_SCOPE, persiste e responde
CONFIRM_COMPOSITE Valida cadastro, prepara resumo/payloads, executa integrações e encaminha

Triagem s=2 com serviço não envia seu texto intermediário. Qualificação s=2 pode produzir texto intermediário, registrado e combinado com a resposta da Coleta. _combine_replies() une os trechos; isso explica por que um turno pode apresentar a conclusão técnica e a solicitação cadastral juntas.

O loop tem teto _MAX_TURN_ITERATIONS=10, independente do limite de Context Packages, para impedir ciclos entre agentes, packages e sessões. Ao excedê-lo, faz handoff por incerteza.

Latest-message-wins e persistência

sequenceDiagram
    participant C as Cliente
    participant W as Webhook
    participant D as DynamoDB
    participant O as Orchestrator
    participant L as LLM
    C->>W: A
    W->>D: latest=A
    O->>L: Executar A
    C->>W: B
    W->>D: latest=B
    L-->>O: Resultado A
    O->>D: Verificar latest=A
    D-->>O: Falha de atualidade
    O->>D: Concluir A como superseded
    O->>L: Executar B com histórico A+B
    L-->>O: Resultado B
    O->>D: Persistir B sob condições
    O-->>C: Resposta B

_ensure_current() testa o guard antes/depois de operações lentas. try_commit() valida latest, pointer ativo e estado. try_put_agent_result() persiste sessão, pointer, sessões substituídas e intenção comercial sob condições atômicas. _abort_if_session_not_agent_owned() relê ownership antes/depois da IA; uma pessoa que assumiu o atendimento não pode ser sobrescrita por resultado atrasado.

Há nova checagem antes do envio e do replay. Se uma resposta persistida perdeu validade, mark_superseded() impede sua reutilização e o invocador exclui esse outbound do histórico. Isso reduz a janela de corrida, mas não pode cancelar uma requisição HTTP já aceita por um provider.

Confirmação e ponto anterior às integrações

execute_confirm_action(), em src/core/application/confirm_action.py, valida nome, e-mail, documento, CEP, UF, telefone, endereço, cidade e descrição. Uma falha incrementa retry_count; na terceira segue para humano por incerteza. Antes de qualquer mutação externa, prepara payloads e resumo. O callback before_integrations verifica novamente o turno e faz o commit lógico.

Em seguida executa hubspot.upsert_contact() e egestor.upsert_customer(). Depois dessa primeira mutação o fluxo termina a sequência iniciada; a persistência final usa require_latest=False para esse caminho, preservando ainda condições de ownership. Não se deve prometer ausência absoluta de efeitos externos de A se B chegou depois desse ponto.

Falha de integração retorna HANDOFF_PENDING com api_failure; sucesso guarda payload e contato. A criação/atualização de deal é responsabilidade de crm_sync após transação. A notificação de handoff ocorre depois de persistir e somente se a mesma sessão segue em HANDOFF_PENDING.

Falhas

Falha Caminho observado
Parse inválido Uma chamada de reparo pelo modelo summary; persistindo erro, handoff llm_invalid_json
Código/enum obrigatório inválido Dispatcher lança erro de aplicação/contrato; handoff visível
IntegrationError no caminho de IA _handle_llm_failure usa llm_timeout inclusive para erros de provider não relacionados a timeout
Package ausente ou teto excedido Handoff incerteza
Notificador retorna falha ou lança exceção Log handoff_notify_failed/handoff_notify_exception; resposta informa fila sem afirmar entrega ao operador
Perda de atualidade ProcessOutcome(superseded=True); nenhuma resposta obsoleta deve ser enviada
Erro não tratado no caso de uso Lambda registra process_immediate_turn_failed e inclui item na lista de falhas SQS

DomainError não está entre as classes capturadas genericamente por process_message(). Uma aresta ilegal que escapar dos guards pode chegar ao retry SQS, em vez de virar automaticamente llm_invalid_json. O método _force_handoff() trata a impossibilidade de mover um estado já terminal com log force_handoff_blocked_by_state.

Testar e alterar

Os testes centrais são tests/unit/application/test_process_message_use_case.py, test_action_mapping.py, test_confirm_action.py, tests/unit/lambdas/test_orchestrator.py e tests/unit/integrations/test_dynamodb_repos.py.

python -m pytest tests/unit/application/test_process_message_use_case.py tests/unit/lambdas/test_orchestrator.py

Ao alterar concorrência, preserve testes de chegada durante IA, durante resumo, entre commit e escrita, takeover humano e replay após falha WhatsApp. Ao adicionar efeitos, defina explicitamente onde entra o guard, qual dado fica durável e se retry pode repetir uma mutação externa. Consulte SQS, DynamoDB e códigos.