Camadas de Memória
em Ambiente Distribuído
Quatro camadas de memória com ciclos de vida completamente diferentes. Redis como checkpointer e working memory centralizado. RAG com OpenSearch por Knowledge Base de agente.
As Quatro Camadas de Memória
O problema central: agentes em contas AWS diferentes, acessados via MCP, não compartilham estado nativo. Cada invocação de tool MCP é stateless por padrão. Você precisa de uma estratégia explícita para cada camada.
Camada 1 — Buffer de Mensagens
O histórico de mensagens vive em state["messages"]. O LangGraph guarda tudo, mas o LLM só deve ver o necessário. Trim sempre — nunca passe o histórico completo para o modelo.
from langchain_core.messages import trim_messages
def supervisor_node(state: AgentState):
mensagens_trimadas = trim_messages(
state["messages"],
strategy="last", # mantém as mensagens mais recentes
token_counter=llm, # conta tokens reais do modelo
max_tokens=8000, # reserva espaço para a resposta
start_on="human", # sempre começa com mensagem humana
include_system=True, # mantém o system prompt
allow_partial=False, # sem mensagens cortadas no meio
)
response = supervisor_llm.invoke(mensagens_trimadas)
return {"messages": [response]}
Camada 2 — Sessão com Checkpointer Redis
O Checkpointer persiste o estado completo do TypedDict entre execuções. É o que permite retomar uma conversa depois que o usuário some por 10 minutos — o grafo carrega o estado anterior automaticamente.
from langgraph.checkpoint.redis import AsyncRedisSaver
import redis.asyncio as aioredis
redis_client = aioredis.from_url(
"redis://elasticache-redis.interno:6379",
decode_responses=True
)
checkpointer = AsyncRedisSaver(redis_client)
# thread_id = identificador único da sessão do usuário
config = {"configurable": {"thread_id": "user_123_session_456"}}
# Invocação 1: usuário faz pergunta
result1 = await app.ainvoke(
{"messages": [HumanMessage("Pesquise sobre X")]},
config
)
# Invocação 2: usuário continua, 5 min depois
# O grafo carrega o estado anterior do Redis automaticamente
result2 = await app.ainvoke(
{"messages": [HumanMessage("Agora analise os dados")]},
config # mesmo thread_id → continua de onde parou
)
Estrutura de Chaves no Redis
| Chave | Conteúdo | TTL |
|---|---|---|
| checkpoint:{thread_id}:{id} | Estado completo serializado do grafo (TypedDict) | 24h |
| checkpoint:{thread_id}:latest | Ponteiro para o último checkpoint | 24h |
| working:{session_id} | dict com resultados acumulados cross-agent | 3h |
| session:user:{user_id}:active | session_id ativo do usuário | 24h |
Camada 3 — Working Memory Cross-Agent
O desafio mais complexo: quando o supervisor invoca agent_pesquisa via MCP e depois invoca agent_analise, como o segundo sabe o que o primeiro encontrou?
Abordagem A — Contexto no payload da tool call
O supervisor lê o Redis e serializa o contexto relevante diretamente no payload da chamada MCP. Simples, auditável, sem dependência nos agentes remotos.
class ContextEnricher:
def __init__(self, redis_client):
self.redis = redis_client
async def build_payload(self, session_id: str, agent_name: str, instrucoes: str) -> dict:
raw = await self.redis.get(f"working:{session_id}")
working_mem = json.loads(raw) if raw else {}
# Filtra contexto relevante por agente — evita payload desnecessário
context_map = {
"agent_pesquisa": ["queries_anteriores", "intent"],
"agent_analise": ["pesquisa_resultado", "intent", "dados_coletados"],
"agent_relatorio": ["pesquisa_resultado", "analise_resultado", "intent"],
}
keys = context_map.get(agent_name, [])
contexto = {k: working_mem[k] for k in keys if k in working_mem}
return {"session_id": session_id, "instrucoes": instrucoes, "contexto": contexto}
async def persist_result(self, session_id: str, agent_name: str, resultado: str):
# Supervisor é o único que escreve no Redis
raw = await self.redis.get(f"working:{session_id}")
working_mem = json.loads(raw) if raw else {}
result_key = agent_name.replace("agent_", "") + "_resultado"
working_mem[result_key] = resultado
await self.redis.setex(f"working:{session_id}", 3600, json.dumps(working_mem))
Abordagem B — Referências ao Redis (para contextos volumosos)
Quando o resultado de um agente é grande demais para o payload, o state guarda apenas a chave Redis, não o dado em si. O próximo agente recebe a chave e busca diretamente (se tiver acesso) ou o supervisor repassa o conteúdo quando necessário.
Camada 4 — RAG com OpenSearch por KB
Cada agente especialista possui seu próprio Knowledge Base no OpenSearch, acessível via VPC Endpoint. Os KBs são completamente independentes — o agente de pesquisa não enxerga o KB de análise e vice-versa.
MCP Central (conta supervisor) │ VPC Endpoint por conta ├──────────────────────────────────────────────────────┐ │ │ │ Conta AWS A Conta AWS B Conta AWS C OpenSearch OpenSearch OpenSearch index: kb-pesquisa index: kb-analise index: kb-relatorio Agente Pesquisa Agente Análise Agente Relatório
Retriever por agente com estratégia diferenciada
from langchain_community.vectorstores import OpenSearchVectorSearch
from langchain_aws import BedrockEmbeddings
def criar_retriever(opensearch_url: str, index: str, aws_region: str, strategy: str):
embeddings = BedrockEmbeddings(
model_id="amazon.titan-embed-text-v2:0",
region_name=aws_region
)
vectorstore = OpenSearchVectorSearch(
opensearch_url=opensearch_url,
index_name=index,
embedding_function=embeddings,
http_auth=get_aws_auth(opensearch_url, aws_region), # SigV4
use_ssl=True,
verify_certs=True,
)
configs = {
"pesquisa": {"search_type": "similarity", "k": 8, "score_threshold": 0.65},
"analise": {"search_type": "mmr", "k": 5, "fetch_k": 20, "lambda_mult": 0.7},
"relatorio": {"search_type": "similarity_score_threshold", "score_threshold": 0.85, "k": 3},
}
return vectorstore.as_retriever(search_kwargs=configs[strategy])
Memória de Longo Prazo Semântica
Diferente do RAG de conhecimento de domínio, a memória semântica guarda fatos sobre o usuário que persistem entre sessões distintas. Implementada como um índice separado no OpenSearch, pesquisável por similaridade.
class LongTermMemoryManager:
async def salvar_fato(self, user_id: str, fato: str):
embedding = await gerar_embedding(fato)
await self.os.index(
index=f"user-memory-{user_id}",
body={
"fato": fato,
"embedding": embedding,
"timestamp": datetime.utcnow().isoformat(),
}
)
async def recuperar_contexto(self, user_id: str, query: str, k: int = 3) -> List[str]:
query_emb = await gerar_embedding(query)
resp = await self.os.search(
index=f"user-memory-{user_id}",
body={"knn": {"embedding": {"vector": query_emb, "k": k}}}
)
return [hit["_source"]["fato"] for hit in resp["hits"]["hits"]]
# Uso no supervisor: enriquece o system prompt com memória do usuário
ltm_fatos = await ltm.recuperar_contexto(
user_id=state["user_id"],
query=state["messages"][-1].content
)
system = f"""Contexto sobre este usuário:
{chr(10).join(ltm_fatos)}
"""
Onde cada tecnologia vive
REDIS (ElastiCache — conta supervisora) ├── Checkpoints LangGraph TTL 24h ├── Working Memory cross-agent TTL 3h └── Cache de sessão de usuário TTL 24h OPENSEARCH (por conta via VPC Endpoint) ├── kb-pesquisa → KB do Agente Pesquisa ├── kb-analise → KB do Agente Análise ├── kb-relatorio → KB do Agente Relatório └── user-memory-* → Memória semântica por usuário EM MEMÓRIA (efêmero — por invocação) └── Buffer de mensagens trimado → direto no AgentState