Ir para o conteúdo principal
Parceiros

Construa agentes duráveis com Temporal e Lakebase

por Sam Ingbar

  • Preserve o progresso do agente com o Temporal: recupere o trabalho registrado após falhas de workers, tente novamente operações com falha e aguarde de forma durável pela revisão humana.
  • Disponibilize o estado da aplicação em tempo real com o Lakebase Postgres: torne evidências, recomendações, decisões de revisão e métricas operacionais consultáveis ao longo de cada execução.
  • Conecte a execução a dados governados: leia as políticas do Unity Catalog por meio de tabelas sincronizadas e, com o Lakebase Change Data Feed ativado, publique alterações operacionais de volta nas tabelas de histórico do Delta.

Um agente de subscrição de empréstimo pessoal coleta evidências, aplica políticas e pode esperar dias por um revisor. Durante esse tempo, os workers podem reiniciar e as chamadas de ferramentas podem falhar. O aplicativo deve preservar o trabalho concluído, retomar a execução e manter as evidências disponíveis para o revisor.

Esta implementação de referência usa o Temporal para execução durável e o Lakebase Postgres para estado operacional consultável. Uma tabela sincronizada disponibiliza a política de subscrição do Unity Catalog no Lakebase. As Atividades do Temporal gravam evidências, decisões e métricas no Lakebase; uma vez ativado, o Lakebase Change Data Feed pode publicar essas alterações em tabelas de histórico Delta gerenciadas pelo Unity Catalog. Essa combinação é especialmente útil quando o Databricks já gerencia as entradas do agente e a análise downstream.

O desafio dos agentes de nuvem de longa duração

Um agente de nuvem pode durar mais que a solicitação, o worker, o contêiner ou a implantação que o iniciou. Um usuário pode iniciar uma sessão, retornar amanhã e continuar em outro worker. Implantações e falhas de processo são rotineiras, portanto, o progresso do agente deve sobreviver independentemente do processo que o executa. A recuperação exige tanto os resultados das operações concluídas quanto o estado do fluxo de controle necessário para determinar o que acontece a seguir.

Para este agente de subscrição, isso cria seis requisitos:

  1. Recuperação: Um worker substituto deve retomar a partir da última etapa concluída.
  2. Tentativas (Retries): Chamadas de ferramentas e operações de banco de dados devem tolerar execuções repetidas sem duplicar efeitos colaterais.
  3. Longas esperas: O agente deve esperar por pessoas ou sistemas externos sem manter um worker aberto.
  4. Visibilidade operacional: Os aplicativos e operadores precisam do status atual, das evidências, do estado de tentativa e dos detalhes da falha.
  5. Governança em tempo de execução: As atualizações de políticas devem ficar disponíveis sem a necessidade de uma implantação de código, e o aplicativo deve definir quando um caso aberto as adota.
  6. Auditoria: O sistema deve reter as evidências, a política, a recomendação e a decisão humana associadas a cada execução.

Uma transcrição de conversa cobre apenas parte desse estado. A recuperação também exige o histórico do fluxo de controle: quais operações foram agendadas, quais resultados foram registrados, pelo que o agente está esperando e quais comandos ele aceitou.

O Temporal simplifica o gerenciamento de sistemas distribuídos. Ao criar com o Temporal, um Workflow é o fluxo de controle durável para uma execução do agente. Uma Activity é uma chamada para um modelo, ferramenta ou banco de dados cujo resultado é registrado no Event History do Workflow; as Activities podem ser repetidas. Um Signal é um comando assíncrono enviado a um Workflow em execução, como a decisão de um subscritor. A implementação de referência do Temporal Lakebase AgentWorkflow é um agente de subscrição de empréstimo pessoal executável. Ele chama várias ferramentas, lê a política governada, gera uma recomendação e aguarda por um subscritor.

O Lakebase Postgres também ajuda os desenvolvedores a gerenciar esses problemas, mas o Temporal e o Lakebase armazenam estados diferentes para consumidores diferentes. O Event History do Temporal impulsiona a reprodução (replay). O Lakebase armazena a visualização voltada para o aplicativo: status de execução atual, mensagens, evidências, estado de revisão e métricas. O Unity Catalog continua sendo a fonte de políticas; uma tabela sincronizada torna essa política consultável no Postgres, e o Change Data Feed fornece o caminho de retorno para o histórico operacional. Os sistemas não compartilham uma transação. As gravações do Lakebase são executadas como Temporal Activities sob execução do tipo "pelo menos uma vez" (at-least-once). Identificadores determinísticos, restrições, atualizações protegidas e upserts do Postgres garantem que as tentativas repetidas de Activity tenham como alvo o mesmo registro lógico.

Essa arquitetura adiciona dois sistemas gerenciados e um contrato de projeção entre eles. Juntos, eles melhoram a resiliência e a escalabilidade do agente, mantendo o custo operacional baixo. O Temporal combinado com o Lakebase é mais útil quando uma sessão de agente precisa sobreviver à substituição de workers, aceitar entradas após longas esperas, expor o estado relacional a um aplicativo e aplicar dados governados enquanto permanece aberta.

O caso de uso de subscrição

Escolhi a subscrição de empréstimos porque a mesma execução deve coletar evidências, aplicar políticas, gerar uma recomendação e esperar por uma pessoa. Um worker pode falhar entre qualquer uma dessas etapas. A política pode mudar sem uma implantação do aplicativo, e a UI precisa das evidências atuais antes que o Workflow seja fechado.

Candidatos simulados (mocked) substituem os órgãos de proteção ao crédito e provedores de renda reais, e a sequência de ferramentas é determinística para fins de simplicidade. Cada solicitação contém um ID de usuário, ID do candidato, valor, finalidade, escolha do modelo e limite de turnos.

O FastAPI atribui run_id, inicia LoanUnderwritingWorkflow e usa esse mesmo ID na API, na execução do Temporal e nas linhas do Lakebase.

No primeiro turno, credit_check retorna a pontuação (score), linhas de crédito, inadimplências e dívida atual. income_verification retorna evidências de renda e emprego. debt_to_income_calc calcula a relação dívida/renda. policy_lookup carrega a política para a finalidade do empréstimo e avalia as evidências em relação aos limites de aprovação, encaminhamento e recusa definitiva (hard-decline).

O candidato limítrofe do exemplo tem uma pontuação de crédito de 665, US$ 76.000 em renda anual verificada, US$ 2.400 em dívidas mensais e um alerta de inadimplência não material. O resultado da política registra cada regra, limite, valor real, resultado de aprovação/reprovação, origem, recomendação e justificativa. O modelo pode recomendar, mas não pode decidir. Um subscritor aprova, nega ou solicita mais informações. Uma solicitação de mais informações torna-se outra mensagem do usuário e outro turno do agente. O caso testa uma falha de worker após chamadas de ferramentas concluídas, uma gravação confirmada no Lakebase cuja conclusão de Activity foi perdida, uma revisão deixada aberta por dias, uma decisão obsoleta do navegador e uma alteração de política durante a execução.

Arquitetura

image1.jpg
Figura 1. Os caminhos de execução, estado operacional e governança na implementação de referência.

Para implementar o agente de subscrição, o React e o FastAPI lidam com o trabalho de HTTP e UI: iniciar execuções, renderizar evidências, listar casos e enviar decisões de revisão. O Temporal Cloud armazena o Event History e despacha as Tasks. Os workers reproduzem o código do Workflow e executam as Activities de modelo, ferramenta e Lakebase; a I/O de rede e de banco de dados permanece fora do código determinístico do Workflow.

Uma execução começa quando o FastAPI inicia um Workflow. O worker agenda as Activities, o Temporal registra seus resultados e o agente eventualmente chega a AWAITING_REVIEW. A resposta do subscritor retorna por meio de um Signal. A aprovação ou negação fecha a execução; uma solicitação de mais informações retoma o loop do agente.

O Lakebase contém dois esquemas operacionais. agent_ops contém o status de execução, mensagens, chamadas de ferramentas, registros de revisão, eventos e métricas que o FastAPI pode consultar com SQL. agent_policy contém a política sincronizada somente leitura usada pelo policy_lookup. Cada Activity grava registros indexados pelos mesmos identificadores determinísticos usados pelo Workflow, para que a projeção possa se atualizar após uma nova tentativa sem tornar o Lakebase parte do mecanismo de replay do Temporal.

O Unity Catalog é a fonte para os limites de subscrição. Uma tabela sincronizada contínua os torna disponíveis para o agente em execução. Os limites aplicados, as evidências e a decisão humana subsequente são gravados em agent_ops. O Change Data Feed pode publicar essas alterações em tabelas de histórico gerenciadas pelo Unity Catalog para auditoria e análise.

Recuperar o trabalho concluído após uma falha de worker

O Temporal mantém o Event History ordenado necessário para reconstruir o estado do Workflow em outro Worker. Esse histórico inclui o agendamento e os resultados de Activity, cronômetros e Signals.

O replay executa o código do Workflow em relação a esses Events registrados e reconstrói variáveis como o turno atual, decisões de revisão aceitas, uso de tokens e evidências coletadas.

Um resultado de Activity registrado é retornado durante o replay em vez de executar a Activity novamente. Uma verificação de crédito concluída permanece concluída, e uma resposta de modelo registrada continua sendo a resposta para essa execução. Se uma Activity estava em andamento quando o Worker falhou e o Temporal nunca registrou sua conclusão, o Temporal pode agendar outra tentativa. Para um agente, isso preserva as respostas do modelo já registradas no Event History. Uma chamada de modelo cuja conclusão não foi registrada ainda pode ser executada novamente, mesmo que o provedor tenha terminado de processá-la.

As Políticas de Repetição (Retry Policies) são atribuídas na granularidade de operações individuais e podem ser reutilizadas no código. No exemplo, as Activities de chamada de modelo permitem até quatro tentativas dentro de um limite de tempo (timeout) de agendamento para fechamento de três minutos. As Activities de chamada de ferramenta permitem até três tentativas e têm um limite de tempo de início para fechamento de 60 segundos. As Activities do Lakebase permitem até cinco tentativas com um limite de tempo de início para fechamento de 15 segundos.

Torne os efeitos externos seguros para repetição

Um dos riscos é que uma gravação de resultado de ferramenta do Lakebase possa ser confirmada (commit) antes que o Worker relate a conclusão da Activity. Se a conexão cair nesse intervalo, o Temporal não terá nenhum resultado registrado e agendará outra tentativa. Ambas as tentativas representam a mesma gravação lógica.

Cada registro do Lakebase possui uma identidade estável. run_id ancora o esquema operacional. message_id identifica uma mensagem, tool_call_id uma invocação de ferramenta, event_id um marco (milestone), review_id uma rodada de revisão e decision_id um comando do revisor. As chaves primárias e restrições de exclusividade (unique constraints) do Postgres aplicam essas identidades.

A gravação de início da ferramenta mostra tanto a identidade estável quanto a proteção de estado terminal (terminal-state guard):

Uma repetição visa o mesmo tool_call_id. O predicado final permite apenas que uma linha não terminal existente seja gravada de volta como iniciada. Se a linha já tiver sido bem-sucedida ou tiver falhado, o PostgreSQL afetará zero linhas. Ele não gera um erro.

O chamador deve inspecionar um resultado de zero linhas. LakebaseWriteResult retorna a contagem de linhas afetadas, mas o wrapper atual da Activity não transforma zero em uma falha. O código de produção deve classificar zero como uma operação nula (no-op) esperada somente após confirmar o estado terminal armazenado; caso contrário, deve gerar ou registrar um conflito. A mesma regra se aplica a transições protegidas de execução e revisão.

Upserts semelhantes cobrem mensagens, resultados de ferramentas e Events. IDs determinísticos fazem com que as repetições converjam para a mesma linha lógica, enquanto cada gravação protegida define quais transições de estado são legais. A API pode exibir brevemente um estado mais antigo enquanto uma gravação é repetida. Depois que a Activity for bem-sucedida, a linha aceita poderá ser consultada.

Toda ferramenta com efeitos colaterais precisa de um contrato equivalente. Uma API de pagamento pode aceitar uma chave de idempotência, um serviço de e-mail um ID de mensagem fornecido pelo chamador e um banco de dados uma restrição de exclusividade. Se o sistema externo não fornecer nenhum mecanismo de eliminação de duplicatas (deduplication), a Activity precisará de seu próprio registro ou de um processo de reconciliação. O Temporal determina quando tentar novamente. O Activity determina como o sistema externo lida com essa repetição.

Expor o estado atual e as métricas operacionais

O Event History fornece semântica de execução e detalhes de depuração. O aplicativo precisa de consultas relacionais indexadas sobre a execução atual: listar casos por usuário e status, carregar uma transcrição com suas evidências, encontrar revisões aguardando por uma pessoa e agregar medições entre as execuções.

O Lakebase armazena essa visualização do aplicativo em um esquema Postgres normalizado. agent_runs mantém o status atual, o Workflow ID, a solicitação, os totais de tokens, os carimbos de data/hora (timestamps) e os metadados de recomendação. agent_messages preserva a transcrição. agent_tool_calls registra argumentos, status, resultado estruturado, erro e tempo. agent_review_decisions conecta a recomendação a um ID de revisão estável, comando do revisor, justificativa e hora da decisão.

O esquema também registra Events nomeados e métricas nos níveis de Workflow, turno e tentativa de Activity. O FastAPI expõe os endpoints run-detail, workflow-metrics e retry-metrics apoiados por essas tabelas. A UI pode mostrar uma execução coletando evidências, outra aguardando revisão e uma terceira repetindo uma ferramenta que falhou. Os operadores podem consultar as mesmas linhas com SQL.

As evidências ficam disponíveis antes da conclusão do Workflow. Depois que policy_lookup terminar, seu resultado estruturado será armazenado com a chamada da ferramenta. Quando a execução atingir AWAITING_REVIEW, o subscritor (underwriter) poderá ver a pontuação de crédito, DTI, limites (thresholds), resultados das regras, justificativa e a origem da política que gerou a recomendação.

Manter a revisão humana durável e rejeitar comandos desatualizados

Quando o modelo retorna uma recomendação, o Workflow deriva review_id de run_id e do turno atual. Ele grava a revisão pendente no Lakebase, registra um evento agent.review_pending, define a projeção como AWAITING_REVIEW e chama workflow.wait_condition. O Temporal retém o Workflow aberto sem manter um processo de Worker ocupado.

A API envia a ação do subscritor como um Signal. Antes de enviá-la, a API verifica se o Lakebase mostra a execução aguardando revisão e se o review_id enviado corresponde à rodada atual. Se qualquer uma das verificações falhar, a API retornará um conflito. O Workflow valida o comando de forma independente em relação ao seu próprio estado e ignora decisões desatualizadas ou duplicadas, protegendo a execução mesmo quando a projeção do Lakebase estiver atrasada.

Após aceitar o Signal, o Workflow persiste a decisão por meio de uma Activity idempotente do Lakebase. A aprovação ou recusa conclui a execução. Uma solicitação de mais informações altera a projeção de volta para RUNNING, anexa a justificativa do revisor como uma mensagem do usuário e inicia o próximo turno. Como o turno mudou, a próxima recomendação recebe um novo review_id.

A resposta 202 da API confirma que o Temporal recebeu o Signal. A aceitação de negócios ocorre de forma assíncrona no Workflow, portanto, um comando pode passar pela pré-verificação da API e ainda assim ser ignorado se o estado de revisão tiver mudado. O cliente atualiza a projeção do Lakebase para observar o estado resultante.

Servir políticas governadas sem reimplantar Workers

Os limites de subscrição mudam independentemente do código do Worker. A tabela de origem no Unity Catalog contém valores específicos para cada finalidade, como pontuação de crédito mínima, DTI de aprovação automática, limites de recusa rígidos e nome da política.

O script de configuração cria uma tabela sincronizada contínua do Lakebase chamada agent_policy.underwriting_policy_limits. policy_lookup que consulta essa cópia somente leitura do Postgres por finalidade de empréstimo normalizada. Os proprietários das políticas atualizam a origem no Unity Catalog; o pipeline de sincronização propaga a alteração e uma execução posterior a lê sem a necessidade de implantação de Worker ou API.

O resultado da política contém os limites aplicados, o valor real de cada regra, o resultado de aprovação/falha e a origem. A demonstração pode recorrer a uma política padrão (fixture policy) quando o Lakebase estiver desativado ou a linha estiver indisponível, e registra esse caminho como fixture_fallback. Um Workflow regulamentado pode, em vez disso, falhar de forma fechada (fail closed). O aplicativo precisa tomar essa decisão de fallback explicitamente.

Retornar alterações operacionais para o Unity Catalog

O repositório prepara cada tabela agent_ops para o Lakebase Change Data Feed definindo REPLICA IDENTITY FULL. Um administrador ainda precisa habilitar o recurso para o esquema. O Lakebase então captura inserções, atualizações e exclusões do log de gravação antecipada (write-ahead log) do Postgres e as grava em lotes em tabelas de histórico do Delta gerenciadas pelo Unity Catalog, nomeadas com o padrão lb_<table>_history.

O Change Data Feed está atualmente em Public Preview e envia (flushes) as alterações aproximadamente a cada 15 segundos. Esse intervalo é adequado para auditoria e análise, enquanto a UI consulta o Lakebase diretamente para obter o estado operacional atual.

As tabelas de histórico podem reconstruir a origem da política de uma execução, evidências de ferramentas, tentativas de Activity, espera de revisão, recomendação e decisão humana. O repositório configura o esquema de origem para esse caminho, mas não inclui uma execução ponta a ponta observada do Change Data Feed. A habilitação do feed e a verificação das tabelas de destino continuam sendo etapas de implantação.

Operar o sistema

A implantação separa o React/FastAPI do Temporal Worker. As réplicas da API escalam com a carga de solicitações; os Workers escalam com o backlog de tarefas de Workflow e Activity e com a simultaneidade configurada. O redimensionamento automático (autoscaling) do Lakebase ajusta a computação do banco de dados dentro dos limites do projeto.

Os preços do Temporal Cloud são baseados em Actions mais o armazenamento ativo e retido do Event History, portanto, a frequência de repetição e históricos abertos por muito tempo também afetam o custo. A equipe ainda define réplicas do Kubernetes, Task Queues, limites de pool de conexões e limites de banco de dados para sua carga de trabalho. Como alternativa, você pode criar seu próprio Temporal Service de código aberto usando a versão de código aberto mais recente.

O cliente do Lakebase usa autenticação máquina a máquina (machine-to-machine) OAuth. Os tokens OAuth do Databricks e as credenciais de banco de dados geradas expiram, portanto, o cliente atualiza seu pool de conexões SQLAlchemy antes que a credencial de banco de dados de uma hora expire. Conexões usam TLS. Sem a rotação, um Worker de longa execução encontraria falhas no banco de dados em um cronograma previsível.

Os operadores usam o Temporal para inspecionar o histórico de Workflow e Activity, o Lakebase para consultar o estado e as métricas do aplicativo, e o Kubernetes para verificar a integridade do processo e da implantação. Um operador pode, então, distinguir uma espera de revisão deliberada de uma repetição de Activity, uma falha de acesso ao banco de dados ou uma ferramenta que falhou.

Evidências e limites

A suíte de testes contém 21 testes aprovados para sequenciamento de Workflow, comportamento de revisão, construção de conexão OAuth, persistência idempotente, contratos de métricas, início de Workflow de API e configurações de Worker. O script de recuperação de falhas adiciona um exercício de falha de processo com o provedor determinístico.

Os dados do candidato e do provedor são fixtures. O repositório não valida modelos de empréstimo, conformidade regulatória, controles de segurança de produção, disponibilidade regional ou desempenho em escala. O exercício de falha local foi executado com o Lakebase desativado, portanto, isola a recuperação do Temporal. O Change Data Feed ainda exige ativação e verificação no ambiente Databricks de destino.

Próximos passos

Quer saber mais? Experimente executar a demonstração você mesmo e veja a execução durável em ação. Execute a implementação de referência, pare um worker no meio da execução e assista ao agente se recuperar. Conecte o Lakebase para explorar as evidências, políticas e o estado de revisão humana por trás de cada decisão.

(Esta publicação no blog foi traduzida utilizando ferramentas baseadas em inteligência artificial) Publicação original

Receba os posts mais recentes na sua caixa de entrada

Assine nosso blog e receba os posts mais recentes diretamente na sua caixa de entrada.