Pesquisar este blog
Automação, desenvolvimento em Python e C#, e soluções práticas para otimizar sua rotina profissional e pessoal. Documentando códigos e aprendizados reais.
Postagem em destaque
- Gerar link
- X
- Outros aplicativos
OpenClaw: agentes para equipes
Problema concreto
Um time de automação precisa orquestrar dezenas de bots que executam tarefas interdependentes em sistemas legados. Cada agente tem seu próprio ciclo de vida, mas falhas de rede, inconsistências de estado e limites de taxa dos serviços externos fazem com que a execução sequencial seja inviável. O desenvolvedor gasta horas debugando deadlocks e reprocessos manuais porque o framework não oferece um modelo de equipe robusto.
Arquitetura de agentes em OpenClaw
Camada de orquestração
OpenClaw expõe Team como um contêiner de Agent. Cada agente roda em seu próprio processo (ou container) e comunica-se via message bus baseado em ZeroMQ. A camada de orquestração cuida de:
- Distribuição de tarefas usando
TaskQueuepersistente (Redis ou PostgreSQL). - Monitoramento de heartbeat para detectar agentes mortos.
- Rebalanceamento automático quando a carga muda.
Essa separação evita que um erro em um agente derrube todo o time, mas introduz latência de rede que deve ser mitigada com batching e backpressure.
Modelo de estado compartilhado
OpenClaw recomenda Event Sourcing para o estado da equipe. Cada mudança gera um evento imutável armazenado em um log (Kafka ou Pulsar). Os agentes podem reconstruir o estado a partir do log, garantindo consistência eventual sem bloqueios.
Implementação de um agente típico
O código abaixo demonstra um agente que processa tickets de suporte. Ele usa asyncio para I/O não bloqueante, type hints e tratamento de exceções granular.
import asyncio
import json
from typing import Any, Dict, Optional
import zmq
import zmq.asyncio
class SupportAgent:
def __init__(self, agent_id: str, broker_url: str) -> None:
self.agent_id = agent_id
self.ctx = zmq.asyncio.Context()
self.socket = self.ctx.socket(zmq.DEALER)
self.socket.identity = agent_id.encode()
self.socket.connect(broker_url)
async def _send_heartbeat(self) -> None:
await self.socket.send_multipart([b'HEARTBEAT', self.agent_id.encode()])
async def _process_message(self, msg: bytes) -> None:
try:
payload: Dict[str, Any] = json.loads(msg)
ticket_id = payload['ticket_id']
# Simula chamada a API externa que pode falhar
await self._handle_ticket(ticket_id)
except (json.JSONDecodeError, KeyError) as exc:
# Log estruturado, não interrompe o loop
print(f"[{self.agent_id}] Mensagem inválida: {exc}")
except Exception as exc:
# Reenfileira a mensagem para retry
await self._requeue(payload)
print(f"[{self.agent_id}] Erro ao processar ticket: {exc}")
async def _handle_ticket(self, ticket_id: str) -> None:
# Exemplo de chamada com timeout e retry exponencial
for attempt in range(3):
try:
await asyncio.wait_for(self._call_external_api(ticket_id), timeout=5)
print(f"[{self.agent_id}] Ticket {ticket_id} concluído")
return
except asyncio.TimeoutError:
await asyncio.sleep(2 ** attempt)
raise RuntimeError(f"Falha ao processar ticket {ticket_id}")
async def _call_external_api(self, ticket_id: str) -> None:
# Placeholder para integração real
await asyncio.sleep(0.1) # Simula latência
async def _requeue(self, payload: Dict[str, Any]) -> None:
# Publica no tópico de retry
await self.socket.send_multipart([b'RETRY', json.dumps(payload).encode()])
async def run(self) -> None:
heartbeat_task = asyncio.create_task(self._heartbeat_loop())
while True:
try:
parts = await self.socket.recv_multipart()
if parts[0] == b'TASK':
await self._process_message(parts[1])
except zmq.error.ZMQError as exc:
print(f"[{self.agent_id}] ZMQ error: {exc}")
break
async def _heartbeat_loop(self) -> None:
while True:
await self._send_heartbeat()
await asyncio.sleep(10)
if __name__ == '__main__':
agent = SupportAgent(agent_id='agent-42', broker_url='tcp://localhost:5555')
asyncio.run(agent.run())
Persistência de eventos
O agente acima publica eventos de retry no broker, mas o estado da equipe deve ser salvo em um log de eventos. Um exemplo mínimo usando confluent-kafka:
from confluent_kafka import Producer, Consumer, KafkaException
producer = Producer({'bootstrap.servers': 'kafka:9092'})
def emit_event(event_type: str, data: dict) -> None:
payload = json.dumps({'type': event_type, 'data': data}).encode()
producer.produce('team-events', payload)
producer.flush()
def consume_events() -> None:
consumer = Consumer({
'bootstrap.servers': 'kafka:9092',
'group.id': 'team-rebuilder',
'auto.offset.reset': 'earliest'
})
consumer.subscribe(['team-events'])
try:
while True:
msg = consumer.poll(1.0)
if msg is None:
continue
if msg.error():
raise KafkaException(msg.error())
event = json.loads(msg.value())
# Reconstruir estado aqui
finally:
consumer.close()
Trade‑offs críticos
Persistência vs. latência
Gravar cada mudança no log garante auditabilidade, mas aumenta a latência de commit. Em fluxos de alta frequência, agrupe eventos em batches de 100 ou use asyncio.Queue para buffer.
Escalabilidade do broker
ZeroMQ escala bem verticalmente, porém não oferece persistência. Quando a equipe ultrapassa 200 agentes, migre para um broker com disco (Kafka, NATS JetStream). O custo de migração inclui refatorar o protocolo de heartbeat para mensagens de confirmação.
Gerenciamento de falhas
Erros de rede são inevitáveis. Estratégia recomendada:
- Retry exponencial com teto de 30 s.
- Circuit breaker por agente (ex.:
pybreaker). - Dead‑letter queue para tickets que falharam três vezes.
Recomendações práticas
- Defina um esquema de evento versionado; mudanças de campo quebram consumidores antigos.
- Use containers leves (Docker) para isolar agentes e facilitar auto‑scaling via Kubernetes Horizontal Pod Autoscaler.
- Monitore métricas de queue depth, heartbeat loss e tempo médio de processamento; alertas devem disparar antes que a fila estoure.
- Implemente testes de integração que simulam falhas de rede e perda de mensagens; use
pytest‑asyncioetoxpara validar.
Armadilhas comuns
- Assumir que o broker é infalível; ignore a necessidade de replay de eventos.
- Compartilhar objetos mutáveis entre agentes sem lock; isso gera race conditions difíceis de reproduzir.
- Manter estado em memória local e confiar em heartbeats para consistência; a perda de um heartbeat não significa que o agente está saudável.
- Gerar link
- X
- Outros aplicativos
Postagens mais visitadas
RPA em Escala: Criando uma Fila de Trabalho para seus Robôs com Python e Redis Queue (RQ)
- Gerar link
- X
- Outros aplicativos
Como preparar dados para aplicações com LLMs
- Gerar link
- X
- Outros aplicativos
Comentários
Postar um comentário