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.