1. Antes de começar
À medida que os agentes de IA lidam com interações longas de várias etapas que abrangem vários dias e executam tarefas de várias etapas de longo prazo, os modelos de linguagem grandes (LLMs) permanecem inerentemente sem estado em todas as sessões. Quando um usuário volta a interagir com um agente no dia seguinte, o modelo começa do zero, a menos que o aplicativo possa reconstruir o contexto necessário.
A abordagem simples para esse problema é o token stuffing, que consiste em anexar históricos de conversas completos, registros de execução de ferramentas e bases de código diretamente a cada solicitação ativa. Embora as janelas de contexto grandes de milhões de tokens tornem isso tecnicamente possível, o preenchimento de contexto introduz um grave problema operacional: os custos de token aumentam quadraticamente a cada turno, as latências de resposta crescem para dezenas de segundos e os modelos sofrem de degradação de contexto "perdido no meio".
Para criar agentes de IA confiáveis, você precisa de uma arquitetura de memória de duas camadas:
- Buffer de sessão de curto prazo: armazena em cache as conversas recentes na memória ativa usando uma janela deslizante limitada por tokens. Esse nível exige pesquisas na memória de alta taxa de transferência e submilissegundos em cada turno, o que torna o Memorystore for Valkey a escolha ideal.
- Memória persistente de longo prazo: armazena entidades estruturadas, preferências do usuário e fatos episódicos em várias sessões. Esse nível exige integridade transacional, segurança multitenant e recuperação híbrida em dados relacionais e vetores, o que torna o AlloyDB para PostgreSQL a escolha certa.

Entenda os quatro tipos de memória
Uma arquitetura de memória robusta depende de quatro tipos de memória complementares ao longo da jornada do usuário:
Tipo de memória | O que ele armazena | Camada de armazenamento | Duração |
Buffer (curto prazo) | Atualizações recentes da conversa bruta | Memorystore for Valkey | Sessão ativa |
Memória de resumo | Histórico compactado de interações mais antigas | Memorystore for Valkey | Janela multiturno |
Memória episódica | Ações, eventos e saídas de ferramentas anteriores | AlloyDB para PostgreSQL (vetor) | Permanente |
Memória de entidades e regras | Preferências, restrições e vetos do usuário | AlloyDB para PostgreSQL (SQL estruturado + vetor) | Permanente |
Impacto medido da memória em camadas
Os testes de benchmark interno em diálogos de desenvolvimento multiturno (mais de 45 rodadas com registros de saída de ferramentas pesados) demonstram uma economia significativa em relação ao preenchimento de contexto simples:
Métrica / dimensão | Context stuffing ingênuo | Memória em camadas (AlloyDB + Memorystore) | Impacto líquido nos testes |
Tamanho do comando ativo (rotação de 45 graus) | 747.033 tokens | 83.262 tokens | Comando 88,9% menor |
Latência de resposta de 45 turnos | 33,5 segundos | 6,7 segundos | Resposta 80% mais rápida |
Tokens de sessão cumulativos | 17,9 milhões de tokens | 4,09 milhões de tokens | 72% de economia total de tokens e custos |
Recall de regras e restrições | Diminui com o tempo | Evita que informações importantes se percam no resumo | Preservado pela pesquisa híbrida |
Atividades deste laboratório
- Provisione o AlloyDB para PostgreSQL e o Memorystore para Valkey.
- Ative o
google_ml_integratione configure as incorporações automáticas transacionais do lado do banco de dados (ai.initialize_embeddings). - Implemente um buffer de sessão do Valkey de curto prazo com um padrão de pipeline de resumo antes do corte.
- Extrair entidades de longo prazo de forma nativa usando as funções de IA nativas do AlloyDB (ou seja,
ai.generate) - Consulte fatos de longo prazo com alta precisão e relevância usando a função de pesquisa híbrida nativa do AlloyDB (
ai.hybrid_search) e a reclassificação da Reciprocal Rank Fusion (RRF). - Crie um avaliador de permissões de ferramentas empresariais de três níveis e um mecanismo de compactação de memória em segundo plano.
- Integre a arquitetura de memória de duas camadas diretamente a um agente autônomo usando o Kit de Desenvolvimento de Agente (ADK) do Google.
O que é necessário
- Ter um projeto do Google Cloud com o faturamento ativado.
- Um navegador da Web, como o Chrome.
- Conhecimento básico de Python e SQL, incluindo experiência com a execução de consultas SQL no AlloyDB (no Studio, na CLI etc.).
Público-alvo e custo
- Público-alvo: desenvolvedores de IA, engenheiros de back-end e arquitetos de banco de dados.
- Custo estimado: os recursos do Google Cloud criados neste codelab custarão aproximadamente US$1,50.
2. Configuração e requisitos
Iniciar o Cloud Shell
Neste codelab, você vai executar comandos no Google Cloud Shell, um terminal hospedado na nuvem pré-configurado com gcloud, psql e python3.
- Abra o Console do Google Cloud.
- Clique em Ativar o Cloud Shell no canto superior direito do console do Google Cloud.
- Verifique a autenticação:
gcloud auth list
export PROJECT_ID=<YOUR_PROJECT_ID>
gcloud config set project $PROJECT_ID
export REGION=us-east1
export GENAI_LOCATION=us
export ZONE=us-east1-b
export ADBCLUSTER=agent-memory-cluster
export ADBINSTANCE=agent-memory-instance
export VALKEYINSTANCE=agent-memory-cache
export VM_NAME=agent-dev-vm
Ativar APIs do Google Cloud e criar uma VM de desenvolvimento
Execute o comando a seguir no Cloud Shell para ativar as APIs necessárias:
gcloud services enable \
alloydb.googleapis.com \
memorystore.googleapis.com \
aiplatform.googleapis.com \
compute.googleapis.com \
servicenetworking.googleapis.com \
networkconnectivity.googleapis.com
Crie uma instância de VM do Compute Engine na rede VPC default para hospedar seu ambiente de desenvolvimento em Python com o AlloyDB e o Memorystore para Valkey:
gcloud compute instances create $VM_NAME \
--zone=$ZONE \
--machine-type=e2-standard-2 \
--scopes=cloud-platform \
--network=default \
--shielded-secure-boot
3. Provisionar o AlloyDB e o Memorystore para Valkey
Nesta etapa, você vai provisionar o cluster e a instância principal do AlloyDB para PostgreSQL, estabelecer uma rede de serviços particulares e criar uma instância do Memorystore para Valkey.
Criar um intervalo de IP de acesso privado a serviços
O AlloyDB exige um intervalo de IP particular na sua rede de nuvem privada virtual (VPC). Supondo que você esteja usando a rede VPC default:
- Crie a alocação de intervalo de IP privado:
gcloud compute addresses create psa-range \
--global \
--purpose=VPC_PEERING \
--prefix-length=24 \
--description="VPC private service access" \
--network=default
- Estabeleça a conexão de peering de VPC particular:
gcloud services vpc-peerings connect \
--service=servicenetworking.googleapis.com \
--ranges=psa-range \
--network=default
Criar cluster e instância principal do AlloyDB
- Crie uma senha inicial do cluster para inicialização do sistema:
export PGPASSWORD=`openssl rand -hex 12`
- Crie um cluster de teste sem custo financeiro:
gcloud alloydb clusters create $ADBCLUSTER \
--password=$PGPASSWORD \
--network=default \
--region=$REGION \
--subscription-type=TRIAL
- Crie a instância principal:
gcloud alloydb instances create $ADBINSTANCE \
--instance-type=PRIMARY \
--cpu-count=2 \
--region=$REGION \
--cluster=$ADBCLUSTER
Provisionar uma instância do Memorystore for Valkey
O Memorystore para Valkey exige uma política de conexão de serviço (gcp-memorystore) na sua rede e região antes da criação da instância.
- Crie a política de conexão de serviço para o Memorystore:
gcloud network-connectivity service-connection-policies create memorystore-policy \
--network=default \
--region=$REGION \
--service-class=gcp-memorystore \
--subnets=projects/$PROJECT_ID/regions/$REGION/subnetworks/default
- Crie a instância do Memorystore para Valkey:
gcloud memorystore instances create $VALKEYINSTANCE \
--location=$REGION \
--shard-count=1 \
--replica-count=0 \
--node-type=SHARED_CORE_NANO \
--psc-auto-connections="network=projects/$PROJECT_ID/global/networks/default,projectId=$PROJECT_ID"
Conceder permissões do IAM da Vertex AI
Conceda à conta de serviço do AlloyDB as permissões do IAM necessárias para invocar os modelos de incorporação da Agent Platform:
PROJECT_ID=$(gcloud config get-value project)
gcloud projects add-iam-policy-binding $PROJECT_ID \
--member="serviceAccount:service-$(gcloud projects describe $PROJECT_ID --format="value(projectNumber)")@gcp-sa-alloydb.iam.gserviceaccount.com" \
--role="roles/aiplatform.user"
4. Inicializar o ambiente e os endpoints de acesso
Configurar a autenticação do IAM do AlloyDB e flags do banco de dados
Ative a autenticação do banco de dados do IAM (alloydb.iam_authentication=on) e o mecanismo de consulta de IA (google_ml_integration.enable_ai_query_engine=on) na sua instância do AlloyDB:
gcloud beta alloydb instances update $ADBINSTANCE \
--cluster=$ADBCLUSTER \
--region=$REGION \
--update-mode=FORCE_APPLY \
--database-flags=alloydb.iam_authentication=on,google_ml_integration.enable_ai_query_engine=on
Em seguida, adicione sua conta do Google Cloud como um usuário do banco de dados baseado no IAM com permissões de superusuário:
export USER_ACCOUNT=$(gcloud config get-value account)
gcloud alloydb users create $USER_ACCOUNT \
--cluster=$ADBCLUSTER \
--region=$REGION \
--type=IAM_BASED \
--db-roles=alloydbsuperuser
Recuperar endpoints de VPC interna no Cloud Shell
Antes de se conectar por SSH à VM de desenvolvimento, recupere os endereços IP internos da VPC para o AlloyDB e o Memorystore para Valkey no Cloud Shell:
# Export GCP Project ID, Region, and IAM User
export PROJECT_ID=$(gcloud config get-value project)
export REGION=us-east1
export DB_USER=$(gcloud config get-value account)
# Retrieve Internal VPC IP Addresses for AlloyDB & Memorystore for Valkey
export DB_HOST=$(gcloud alloydb instances describe $ADBINSTANCE --cluster=$ADBCLUSTER --region=$REGION --format="value(ipAddress)")
export DB_PORT=5432
export DB_NAME=postgres
export VALKEY_HOST=$(gcloud memorystore instances describe $VALKEYINSTANCE --location=$REGION --format="value(discoveryEndpoints[0].address)")
export VALKEY_PORT=6379
# Verify exported endpoints
echo "Project ID: $PROJECT_ID"
echo "Region: $REGION"
echo "AlloyDB Private IP: $DB_HOST"
echo "IAM DB User: $DB_USER"
echo "Valkey Private IP: $VALKEY_HOST"
Acessar a VM de desenvolvimento via SSH e exportar variáveis de conexão
Faça SSH do Cloud Shell para a VM de desenvolvimento do Compute Engine (agent-dev-vm) localizada na mesma rede VPC:
gcloud compute ssh $VM_NAME --zone=$ZONE
Depois de fazer login na VM de desenvolvimento, exporte a configuração do projeto e os endpoints de conexão acima. Substitua pelo endereço de e-mail exato usado ao criar o usuário do AlloyDB:
export PROJECT_ID=<YOUR_PROJECT_ID>
export REGION=us-east1
export GENAI_LOCATION=us
export DB_HOST=<YOUR_ALLOYDB_PRIVATE_IP>
export DB_PORT=5432
export DB_NAME=postgres
export DB_USER=<YOUR_IAM_DB_USER>
export VALKEY_HOST=<YOUR_VALKEY_PRIVATE_IP>
export VALKEY_PORT=6379
Inicializar o ambiente virtual do Python
Na VM de desenvolvimento, primeiro crie o diretório de trabalho local:
mkdir -p ~/alloydb_agent_memory && cd ~/alloydb_agent_memory
Agora, vamos instalar os pacotes do ambiente virtual Python do sistema, autenticar o Application Default Credentials (ADC) e configurar seu espaço de trabalho:
sudo apt-get update && sudo apt-get install -y python3-venv python3-pip
python3 -m venv venv
source venv/bin/activate
Por fim, no novo ambiente virtual, vamos instalar as dependências:
gcloud auth application-default login
pip install psycopg2-binary valkey redis google-genai
5. Arquitetura do sistema e hierarquia do módulo de código
Visão geral do que você está criando
Antes de implementar os scripts individuais em Python, revise a arquitetura do sistema abaixo. Este aplicativo de exemplo é estruturado em sete scripts modulares do Python que operam em dois caminhos de execução principais que interagem com a mesma instância de banco de dados do AlloyDB para PostgreSQL:
- Caminho de leitura (
hybrid_retriever.py): decompõe perguntas complexas de várias partes em subconsultas de aspecto único diretamente no PostgreSQL usando oai.generate()da IA do AlloyDB e consulta memórias de longo prazo usando a pesquisa híbrida nativa do AlloyDB (ai.hybrid_search). - Write Path (
async_worker.py): um worker de fila em segundo plano fora da linha de execução que extrai de forma assíncrona fatos de entidades estruturadas de trocas de diálogo usando o Gemini Flash e os insere ou atualiza noagent_entities.

Hierarquia de módulos e funções do sistema
Arquivo do módulo | Camada do sistema | Responsabilidade principal |
| Camada de conexão | Estabelece a autenticação do IAM criptografada com SSL no AlloyDB e conexões resilientes a soquetes na Memorystore para Valkey. |
| Memória de curto prazo | Gerencia o histórico de sessões com menos de um milésimo de segundo no Valkey, implementando resumos rotativos Summarize-Before-Trim. |
| Gravador de caminho de trabalho | Executa um worker de fila em segundo plano de daemon off-thread que extrai fatos de entidades usando o Gemini Flash e os insere ou atualiza no AlloyDB. |
| Read Path Retriever | Decompõe perguntas complexas em subconsultas de aspecto único usando a IA do AlloyDB |
| Loop principal do agente | Coordena o loop de execução de turnos de ponta a ponta: busca de curto prazo, pesquisa de longo prazo, montagem de comandos, execução de LLM e enfileiramento assíncrono. |
| Governança e administração | Aplica políticas de segurança de três níveis para execução de ferramentas e agrega memória histórica. |
| Teste e avaliação | Suíte de verificação principal que executa cenários multiturno e de várias sessões, mede a porcentagem de economia de tokens e verifica a precisão da memória. |
6. Configurar o esquema da IA do AlloyDB e os embeddings automáticos transacionais
Visão geral do objetivo e da arquitetura
Neste módulo, você vai definir o esquema do banco de dados do AlloyDB para memória episódica e de longo prazo, estratégias de indexação e autoembedding do lado do banco de dados.
- Repositório de vetores episódicos (
episodic_memory_embeddings): partes não estruturadas de transcrições de chat indexadas com índices de vetores HNSW (vector_cosine_ops). - Repositório de entidades de longo prazo (
agent_entities): fatos estruturados, escolhas do usuário e regras do projeto armazenados com metadados de escopo (global,project,session). Inclui uma coluna de pesquisa de texto completo do PostgreSQL gerada automaticamente (summary_tsv) indexada via RUM. - Auto-Embeddings do lado do Database (
ai.initialize_embeddings): incorpora automaticamente linhas de texto simples novas ou atualizadas nosummary_embeddingusando a Agent Platformtext-embedding-005em segundo plano.
Conectar-se ao AlloyDB Studio
- Acesse a página AlloyDB para Postgres no console do Google Cloud.
- Clique na instância principal.
- Na navegação à esquerda, clique em AlloyDB Studio.
- Selecione o banco de dados
postgres. - Autenticar com
IAM database authentication
Implementação e código-fonte
Depois de se conectar ao banco de dados do AlloyDB PostgreSQL, execute as consultas DDL abaixo:
-- 1. Enable google_ml_integration extension
CREATE EXTENSION IF NOT EXISTS google_ml_integration CASCADE;
CREATE EXTENSION IF NOT EXISTS vector CASCADE;
CREATE EXTENSION IF NOT EXISTS rum CASCADE;
SET google_ml_integration.enable_preview_ai_functions = true;
-- 2. Create Episodic Memory Vector Store Tables
CREATE TABLE IF NOT EXISTS episodic_memory_collections (
uuid UUID PRIMARY KEY,
name VARCHAR,
cmetadata JSONB
);
CREATE TABLE IF NOT EXISTS episodic_memory_embeddings (
uuid UUID PRIMARY KEY,
collection_id UUID REFERENCES episodic_memory_collections(uuid) ON DELETE CASCADE,
embedding VECTOR(768),
document VARCHAR,
cmetadata JSONB,
custom_id VARCHAR,
created_at TIMESTAMPTZ DEFAULT NOW()
);
CREATE INDEX IF NOT EXISTS episodic_memory_embedding_idx
ON episodic_memory_embeddings
USING hnsw (embedding vector_cosine_ops)
WITH (m = 16, ef_construction = 64);
-- 3. Create Entity & Preference Table with Metadata Scope & Auto-Generated TSVector Column
CREATE TABLE IF NOT EXISTS agent_entities (
entity_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
user_id TEXT NOT NULL,
entity_name TEXT NOT NULL,
project_id TEXT DEFAULT 'global',
session_id TEXT,
scope TEXT NOT NULL DEFAULT 'project', -- 'global', 'project', 'session'
summary TEXT NOT NULL,
summary_embedding VECTOR(768),
summary_tsv TSVECTOR GENERATED ALWAYS AS (
to_tsvector('english', entity_name || ' ' || summary)
) STORED,
updated_at TIMESTAMPTZ DEFAULT NOW(),
UNIQUE (user_id, entity_name)
);
CREATE INDEX IF NOT EXISTS agent_entities_scope_idx
ON agent_entities (user_id, project_id, scope);
CREATE INDEX IF NOT EXISTS agent_entities_embedding_idx
ON agent_entities USING hnsw (summary_embedding vector_cosine_ops);
CREATE INDEX IF NOT EXISTS agent_entities_tsv_idx
ON agent_entities USING rum (summary_tsv rum_tsvector_ops);
-- 4. Create User Permissions Table for Tool Execution Controls
CREATE TABLE IF NOT EXISTS user_permissions (
permission_id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
user_id TEXT NOT NULL,
project_id TEXT NOT NULL DEFAULT 'global',
tool_name TEXT NOT NULL,
command_pattern TEXT NOT NULL,
action TEXT NOT NULL CHECK (action IN ('ALLOW', 'BLOCK')),
created_at TIMESTAMPTZ DEFAULT NOW(),
updated_at TIMESTAMPTZ DEFAULT NOW()
);
CREATE INDEX IF NOT EXISTS idx_user_permissions_lookup
ON user_permissions (user_id, project_id, tool_name);
-- 5. Register Gemini 3.5 Flash Model Endpoint in AlloyDB
-- Replace PROJECT_ID with your Google Cloud project ID
CALL google_ml.create_model(
model_id => 'gemini-3.5-flash',
model_request_url => 'https://aiplatform.googleapis.com/v1/projects/PROJECT_ID/locations/global/publishers/google/models/gemini-3.5-flash:generateContent',
model_qualified_name => 'gemini-3.5-flash',
model_provider => 'google',
model_type => 'llm',
model_auth_type => 'alloydb_service_agent_iam'
);
Inicializar incorporações automáticas transacionais e registrar o modelo do Gemini
Em seguida, execute as instruções CALL em blocos separados de execução de consultas para registrar o processo em segundo plano de incorporação automática e o endpoint do modelo Gemini 3.5 Flash:
-- 6. Register Database-Side Transactional Auto-Embedding
CALL ai.initialize_embeddings(
model_id => 'text-embedding-005',
table_name => 'agent_entities',
content_column => 'summary',
embedding_column => 'summary_embedding',
incremental_refresh_mode => 'transactional',
batch_size => 10
);
7. Configurar clientes de conexão do AlloyDB e do Valkey
Visão geral do objetivo e da arquitetura
Neste módulo, você vai estabelecer conexões de rede seguras com o AlloyDB para PostgreSQL (memória de longo prazo) e o Memorystore para Valkey (cache de curto prazo).
- Autenticação do IAM do AlloyDB: usa
gcloud auth application-default print-access-tokenpara recuperar um token OAuth2 de curta duração para conexões de banco de dados sem senha e criptografadas com SSL (sslmode="require"). - Resiliência de rede do Valkey: configura o
redis.Rediscom um tempo limite de soquete de 5,0 segundos (socket_timeout=5.0) para processar operações de rede VPC com segurança em instâncias do Valkey de nó único ou em cluster.
Implementação e código-fonte
Crie o script db_clients.py no diretório de trabalho:
import os
import subprocess
import psycopg2
import redis
# Initialize database & Valkey connection clients from environment variables
DB_HOST = os.getenv("DB_HOST", "127.0.0.1")
DB_PORT = os.getenv("DB_PORT", "5432")
DB_NAME = os.getenv("DB_NAME", "postgres")
DB_USER = os.getenv("DB_USER")
VALKEY_HOST = os.getenv("VALKEY_HOST", "127.0.0.1")
VALKEY_PORT = int(os.getenv("VALKEY_PORT", "6379"))
def get_db_connection():
# Fetch Application Default Credentials (ADC) token for IAM Database Authentication & enable SSL encryption
access_token = subprocess.check_output(
["gcloud", "auth", "application-default", "print-access-token"], text=True
).strip()
return psycopg2.connect(
host=DB_HOST,
port=DB_PORT,
dbname=DB_NAME,
user=DB_USER,
password=access_token,
sslmode="require"
)
def get_valkey_client():
return redis.Redis(
host=VALKEY_HOST,
port=VALKEY_PORT,
db=0,
socket_timeout=5.0,
socket_connect_timeout=5.0
)
if __name__ == "__main__":
print(f"Connecting to AlloyDB via IAM Auth ({DB_USER}) at {DB_HOST}:{DB_PORT} and Valkey at {VALKEY_HOST}:{VALKEY_PORT}...")
print("Database and Valkey connection client modules loaded successfully.")
8. Armazenar em cache o estado da sessão de curto prazo no Memorystore for Valkey
Visão geral do objetivo e da arquitetura
Neste módulo, você vai criar um cache de contexto de curto prazo com menos de um milissegundo no Memorystore para Valkey, implementando um padrão Resumir antes de cortar automatizado.
- Janela deslizante do Valkey: as conversas ativas são armazenadas como strings JSON na chave
session:{session_id}:turns. - Tags de hash do Redis (
{session_id}): a formatação de chavessession:{session_id}:turnsesession:{session_id}:summaryusa tags de hash do cluster do Redis ({...}), forçando as duas chaves ao mesmo slot de hash para garantir a execução atômica em qualquer implantação do Valkey de nó único ou em cluster. - Summarize-Before-Trim: quando a contagem de turnos excede
trigger_limit, os turnos mais antigos que serão cortados são resumidos pelo Gemini Flash em um resumo de texto contínuo (session:{session_id}:summary) antes de reduzir o histórico bruto parawindow_size.
Implementação e código-fonte
Crie o script valkey_buffer.py no diretório de trabalho:
import json
import logging
import time
from typing import Any, Dict, List
import redis
logger = logging.getLogger(__name__)
IN_MEMORY_VALKEY_FALLBACK: Dict[str, Any] = {}
VALKEY_COMPACTION_TOKENS = 0
def get_valkey_compaction_tokens() -> int:
return VALKEY_COMPACTION_TOKENS
def append_session_turn_with_rolling_summary(
valkey_client: redis.Redis,
llm_client: Any,
session_id: str,
user_msg: str,
ai_msg: str,
trigger_limit: int = 10,
window_size: int = 4
) -> None:
"""Appends turn to Valkey. Before trimming old turns, summarizes them into a rolling summary."""
global VALKEY_COMPACTION_TOKENS
# Use Redis Hash Tags {session_id} so both keys hash to the same cluster slot
turns_key = "session:{" + session_id + "}:turns"
summary_key = "session:{" + session_id + "}:summary"
try:
# 1. Append new turn messages
valkey_client.rpush(turns_key, json.dumps({"role": "user", "content": user_msg}))
valkey_client.rpush(turns_key, json.dumps({"role": "assistant", "content": ai_msg}))
raw_turns = valkey_client.lrange(turns_key, 0, -1)
# 2. Check if total turns exceed summary trigger threshold
msg_trigger_count = trigger_limit * 2
msg_window_count = window_size * 2
if len(raw_turns) > msg_trigger_count:
turns_to_prune = raw_turns[:-msg_window_count]
existing_summary = valkey_client.get(summary_key)
existing_summary_text = (
existing_summary.decode('utf-8') if isinstance(existing_summary, bytes) else (existing_summary or "")
)
pruned_text = "\n".join([
f"{json.loads(t)['role']}: {json.loads(t)['content']}" for t in turns_to_prune
])
prompt = f"""Update the existing rolling conversation summary with the old turns about to be trimmed.
Retain key constraints, user choices, technical decisions, and ongoing goals.
Existing Summary:
{existing_summary_text if existing_summary_text else 'No prior summary.'}
Old turns about to be trimmed:
{pruned_text}
Updated Rolling Summary:"""
VALKEY_COMPACTION_TOKENS += max(1, len(prompt) // 4)
updated_summary_text = llm_client.models.generate_content(
model="gemini-3.5-flash", contents=prompt
).text.strip()
# Execute commands directly to support all Redis/Valkey cluster topologies
valkey_client.set(summary_key, updated_summary_text)
valkey_client.ltrim(turns_key, -msg_window_count, -1)
logger.info("Updated rolling summary and trimmed Valkey buffer for session %s", session_id)
except (redis.exceptions.RedisError, Exception) as e:
logger.warning("Valkey operation warning (%s). Falling back to in-memory short-term buffer.", e)
if turns_key not in IN_MEMORY_VALKEY_FALLBACK:
IN_MEMORY_VALKEY_FALLBACK[turns_key] = []
IN_MEMORY_VALKEY_FALLBACK[turns_key].append({"role": "user", "content": user_msg})
IN_MEMORY_VALKEY_FALLBACK[turns_key].append({"role": "assistant", "content": ai_msg})
# In-memory summarize-before-trim fallback logic
msg_trigger_count = trigger_limit * 2
msg_window_count = window_size * 2
raw_fallback_turns = IN_MEMORY_VALKEY_FALLBACK[turns_key]
if len(raw_fallback_turns) > msg_trigger_count:
turns_to_prune = raw_fallback_turns[:-msg_window_count]
existing_summary_text = IN_MEMORY_VALKEY_FALLBACK.get(summary_key, "")
pruned_text = "\n".join([f"{t['role']}: {t['content']}" for t in turns_to_prune])
prompt = f"""Update the existing rolling conversation summary with the old turns about to be trimmed.
Retain key constraints, user choices, technical decisions, and ongoing goals.
Existing Summary:
{existing_summary_text if existing_summary_text else 'No prior summary.'}
Old turns about to be trimmed:
{pruned_text}
Updated Rolling Summary:"""
VALKEY_COMPACTION_TOKENS += max(1, len(prompt) // 4)
for attempt in range(4):
try:
updated_summary_text = llm_client.models.generate_content(
model="gemini-3.5-flash", contents=prompt
).text.strip()
break
except Exception as e:
if ("429" in str(e) or "RESOURCE_EXHAUSTED" in str(e)) and attempt < 3:
time.sleep(3 * (2 ** attempt))
else:
raise
IN_MEMORY_VALKEY_FALLBACK[summary_key] = updated_summary_text
IN_MEMORY_VALKEY_FALLBACK[turns_key] = raw_fallback_turns[-msg_window_count:]
def get_session_context_buffer(
valkey_client: redis.Redis,
session_id: str
) -> Dict[str, Any]:
"""Retrieves rolling summary + sliding window history from Valkey to build prompt context."""
turns_key = "session:{" + session_id + "}:turns"
summary_key = "session:{" + session_id + "}:summary"
try:
summary = valkey_client.get(summary_key)
summary_text = summary.decode('utf-8') if isinstance(summary, bytes) else (summary or "")
raw_turns = valkey_client.lrange(turns_key, 0, -1)
recent_turns = [json.loads(t) for t in raw_turns]
return {
"rolling_summary": summary_text,
"recent_turns": recent_turns
}
except (redis.exceptions.RedisError, Exception) as e:
logger.warning("Valkey read error (%s). Using in-memory short-term fallback buffer.", e)
summary_text = IN_MEMORY_VALKEY_FALLBACK.get(summary_key, "")
recent_turns = IN_MEMORY_VALKEY_FALLBACK.get(turns_key, [])
return {"rolling_summary": summary_text, "recent_turns": recent_turns}
9. Extrair entidades fora da linha de execução (worker de memória em segundo plano)
Visão geral do objetivo e da arquitetura
Neste módulo, você vai criar um worker de extração em segundo plano do caminho de gravação off-thread (AsyncMemoryWorker) que extrai fatos de entidades de longo prazo sem diminuir a velocidade das respostas interativas de IA.
- Trabalhador de fila não bloqueadora: inicia uma linha de execução de daemon (
queue.Queue) para que as conversas com o desenvolvedor retornem imediatamente sem esperar a extração do LLM ou as gravações no banco de dados. - Extração de fatos de entidades fora da linha de execução: chama o Gemini Flash (
model="gemini-3.5-flash",response_mime_type="application/json") em segundo plano para analisar entidades estruturadas sem bloquear as conversas com o usuário. - Resolução de correferência temporal (
build_temporal_rules_prompt): aplica regras que convertem expressões temporais relativas (por exemplo, "atualmente", "última sessão") em IDs de sessão explícitos (por exemplo,session_id). - Inserções/atualizações de esquema: solicita ao Gemini Flash que retorne uma matriz JSON de entidades (
entity_name,project_id,scope,summary) e grava texto simples emagent_entitiesusando instruçõesON CONFLICT DO UPDATEparametrizadas do PostgreSQL.
Implementação e código-fonte
Crie o script async_worker.py no diretório de trabalho:
import json
import logging
import queue
import threading
import time
from typing import Dict, Any, Optional
import psycopg2
logger = logging.getLogger(__name__)
def build_temporal_rules_prompt(active_session_id: str, prev_session_id: Optional[str] = None) -> str:
"""REUSABLE PROMPT HELPER: Defines temporal coreference rules for both Read and Write paths."""
return f"""TEMPORAL COREFERENCE RULES:
1. Do not use relative temporal words like 'currently', 'now', 'this session', 'previous session', or 'last session'.
2. Replace any relative temporal reference with explicit session identifiers:
- Active Session: '{active_session_id}'
- Previous Session: '{prev_session_id if prev_session_id else active_session_id}'
(e.g., write 'User prefers Python in session {active_session_id}' instead of 'User currently prefers Python')."""
class AsyncMemoryWorker:
"""Off-thread daemon queue worker performing background entity extraction and upserts into AlloyDB."""
def __init__(self, conn_factory, llm_client):
self.conn_factory = conn_factory
self.llm_client = llm_client
self.work_queue = queue.Queue()
self.total_extraction_tokens = 0
self.worker_thread = threading.Thread(target=self._run_worker, daemon=True)
self.worker_thread.start()
def enqueue_extraction(
self, user_id: str, session_id: str, user_msg: str, ai_msg: str, prev_session_id: Optional[str] = None
):
self.work_queue.put((user_id, session_id, user_msg, ai_msg, prev_session_id))
def _run_worker(self):
while True:
item = self.work_queue.get()
if item is None:
break
user_id, session_id, user_msg, ai_msg, prev_session_id = item
try:
self._extract_and_upsert(user_id, session_id, user_msg, ai_msg, prev_session_id)
except Exception as e:
logger.error("Background extraction failed: %s", e)
finally:
self.work_queue.task_done()
def _extract_and_upsert(
self, user_id: str, session_id: str, user_msg: str, ai_msg: str, prev_session_id: Optional[str] = None
):
temporal_rules = build_temporal_rules_prompt(session_id, prev_session_id)
json_fmt = '[{"entity_name": "...", "project_id": "...", "scope": "global|project|session", "summary": "..."}]'
prompt = f"""Extract key preferences and named entities from this exchange.
Format output strictly as a JSON array of objects with structure: {json_fmt}.
{temporal_rules}
Rules for Scope:
- 'global': Universal developer preferences.
- 'project': Architecture, database choices, timeouts, and rules for a specific project.
- 'session': Ephemeral task state tied strictly to a single session.
User: {user_msg}
AI: {ai_msg}
Output (JSON array only):"""
self.total_extraction_tokens += max(1, len(prompt) // 4)
for attempt in range(4):
try:
response = self.llm_client.models.generate_content(
model="gemini-3.5-flash",
contents=prompt,
config={"response_mime_type": "application/json"}
)
break
except Exception as e:
if ("429" in str(e) or "RESOURCE_EXHAUSTED" in str(e)) and attempt < 3:
time.sleep(3 * (2 ** attempt))
else:
raise
content = str(response.text).strip()
if not content:
return
extracted_list = json.loads(content)
if not isinstance(extracted_list, list):
return
upsert_sql = """
INSERT INTO agent_entities (user_id, entity_name, project_id, session_id, scope, summary, updated_at)
VALUES (%s, %s, %s, %s, %s, %s, NOW())
ON CONFLICT (user_id, entity_name)
DO UPDATE SET
project_id = EXCLUDED.project_id,
session_id = EXCLUDED.session_id,
scope = EXCLUDED.scope,
summary = EXCLUDED.summary,
updated_at = NOW();
"""
conn = self.conn_factory()
try:
with conn:
with conn.cursor() as cur:
for item in extracted_list:
e_name = item.get("entity_name", "General_Fact")
p_id = item.get("project_id", "global")
s_scope = item.get("scope", "project")
summary = item.get("summary", "")
if e_name and summary:
cur.execute(upsert_sql, (user_id, e_name, p_id, session_id, s_scope, summary))
logger.info("Saved %d entities for user %s", len(extracted_list), user_id)
finally:
conn.close()
10. Consultar a memória de longo prazo com a pesquisa híbrida nativa
Visão geral do objetivo e da arquitetura
Neste módulo, você vai implementar a decomposição de subconsultas do caminho de leitura, a reescrita de consultas temporais, o isolamento do escopo de metadados e usar a pesquisa híbrida nativa do AlloyDB (ai.hybrid_search).
- Decomposição de subconsultas no banco de dados (
rewrite_and_decompose_query): usa a função integradaai.generate()da AlloyDB AI diretamente no PostgreSQL para dividir perguntas complexas em subconsultas de aspecto único com IDs de sessão normalizados, evitando que consultas de vários tópicos reduzam a precisão da pesquisa de vetores. - Filtragem de escopo no nível do índice: cria filtros SQL do lado do banco de dados (
filter_condition => "user_id = '%s' AND (project_id = '%s' OR scope = 'global')") para isolar memórias por usuário e projeto, incluindo preferências globais do desenvolvedor. - Pesquisa híbrida nativa do AlloyDB (
ai.hybrid_search): combina a similaridade de cosseno vetorial (public.<=>) com a pesquisa de texto completo (rum) no AlloyDB usando a Reciprocal Rank Fusion (RRF) para oferecer acurácia, recall e relevância ideais.
Implementação e código-fonte
Crie o script hybrid_retriever.py no diretório de trabalho:
import os
import json
import logging
from typing import Any, Dict, List, Optional
import psycopg2
from psycopg2.extras import RealDictCursor
from async_worker import build_temporal_rules_prompt
logger = logging.getLogger(__name__)
def rewrite_and_decompose_query(
conn: Any,
llm_client: Any,
text: str,
active_session_id: str,
prev_session_id: Optional[str] = None
) -> List[str]:
"""Read-Path LLM Query Normalizer & Sub-Query Decomposer using AlloyDB AI ai.generate() directly in PostgreSQL."""
temporal_rules = build_temporal_rules_prompt(active_session_id, prev_session_id)
prompt = f"""You are a query normalization and decomposition tool for an AI agent's memory retrieval system.
Tasks:
1. Normalize any relative temporal phrases in the user query into explicit session identifiers.
{temporal_rules}
2. Decompose compound user queries asking about multiple distinct topics into up to 4 concise, single-topic search queries. Ensure all distinct questions (both general developer preferences and project-specific architecture/rules) are preserved as separate sub-queries.
3. Return ONLY a valid JSON array of strings containing the sub-queries.
Example Output Format:
["general developer coding preferences", "CloudRetail backend stack in session sess_2026_01", "CloudRetail timeout rules in session sess_2026_01"]
Text to Process: {text}
JSON Output:"""
try:
with conn.cursor() as cur:
cur.execute("SET google_ml_integration.enable_preview_ai_functions = true;")
cur.execute("SELECT ai.generate(prompt => %s::text, model_id => 'gemini-3.5-flash'::varchar(100));", (prompt,))
row = cur.fetchone()
if row and row[0]:
res_text = str(row[0]).strip()
clean_text = res_text
if clean_text.startswith("```"):
clean_text = clean_text.removeprefix("```json").removeprefix("```").removesuffix("```").strip()
parsed = json.loads(clean_text)
if isinstance(parsed, list) and len(parsed) > 0:
return [str(q).strip() for q in parsed]
return [text]
except Exception as e:
logger.error("Error in in-database ai.generate sub-query decomposition: %s", e)
return [text]
def rewrite_temporal_query(
conn: Any,
llm_client: Any,
text: str,
active_session_id: str,
prev_session_id: Optional[str] = None
) -> str:
"""Backward-compatible wrapper returning first decomposed query string."""
queries = rewrite_and_decompose_query(conn, llm_client, text, active_session_id, prev_session_id)
return queries[0] if queries else text
def retrieve_hybrid_entities(
conn,
llm_client: Any,
user_id: str,
active_session_id: str,
query_text: str,
query_embedding: Optional[List[float]] = None,
prev_session_id: Optional[str] = None,
project_id: Optional[str] = None,
k: int = 20
) -> Dict[str, Any]:
"""Queries long-term entities using Sub-Query Decomposition, metadata scope filtering, and parameterized AlloyDB hybrid search."""
sub_queries = rewrite_and_decompose_query(
conn=conn, llm_client=llm_client, text=query_text, active_session_id=active_session_id, prev_session_id=prev_session_id
)
all_search_queries = list(dict.fromkeys(sub_queries + [query_text]))
retrieved_entities: Dict[str, Any] = {}
filter_cond = f"user_id = '{user_id}' AND (project_id = '{project_id}' OR scope = 'global')" if project_id else f"user_id = '{user_id}'"
query_sql = """
SELECT e.entity_name, e.summary, e.project_id, e.scope, e.updated_at
FROM ai.hybrid_search(
search_inputs => %s::JSONB[],
include_json_output => false
) h
JOIN agent_entities e ON e.entity_name = h.id
WHERE e.user_id = %s
LIMIT %s;
"""
from google import genai
embed_client = genai.Client(
vertexai=True,
project=os.getenv("PROJECT_ID"),
location=os.getenv("REGION", "us-east1")
)
for sq in all_search_queries:
try:
emb_resp = embed_client.models.embed_content(model="text-embedding-005", contents=sq)
sq_embedding = emb_resp.embeddings[0].values
vector_literal = f"'{sq_embedding}'::vector"
except Exception as e:
logger.error(f"Failed to generate embedding for sub-query: {sq}. Error: {e}")
# Fallback to zero vector to prevent Postgres from crashing
sq_embedding = [0.0] * 768
vector_literal = f"'{sq_embedding}'::vector"
search_inputs = [
json.dumps({
"data_type": "vector",
"weight": 0.4,
"table_name": "agent_entities",
"key_column": "entity_name",
"vec_column": "summary_embedding",
"distance_operator": "public.<=>",
"limit": 20,
"query_vector": vector_literal,
"filter_condition": filter_cond
}),
json.dumps({
"data_type": "text",
"weight": 0.6,
"table_name": "agent_entities",
"key_column": "entity_name",
"text_column": "summary_tsv",
"limit": 20,
"ranking_function": "ts_rank",
"query_text_input": sq,
"filter_condition": filter_cond
})
]
with conn.cursor(cursor_factory=RealDictCursor) as cur:
try:
cur.execute("SET google_ml_integration.enable_preview_ai_functions = true;")
cur.execute(query_sql, (search_inputs, user_id, k))
rows = cur.fetchall()
for row in rows:
name = row["entity_name"]
if name not in retrieved_entities:
retrieved_entities[name] = {
"summary": row["summary"],
"project_id": row["project_id"],
"scope": row["scope"],
"updated_at": row["updated_at"].isoformat() if row["updated_at"] else None
}
except Exception as e:
logger.error("Error executing hybrid search for sub-query '%s': %s", sq, e)
return retrieved_entities
11. Criar e executar o loop de memória do agente de ponta a ponta
Visão geral do objetivo e da arquitetura
Neste módulo, você vai criar a principal função de orquestração de agentes (run_agent_turn) que combina recuperação de cache de curto prazo, pesquisa de memória de longo prazo, montagem de comandos, geração de LLM e extração de memória em segundo plano.
- Contexto de curto prazo: busca turnos de diálogo ativos do Valkey e um resumo contínuo (
get_session_context_buffer). - Pesquisa de longo prazo: consulta o AlloyDB via
retrieve_hybrid_entitiesusando subconsultas decompostas filtradas por ID do projeto ativo e escopo global. - Formatação do comando do sistema: reúne
build_agent_promptcom entidades de longo prazo, resumo de curto prazo, diálogo recente e o comando do usuário em um comando do sistema eficiente em termos de tokens. - Enfileiramento assíncrono: armazena em cache a vez no Valkey e enfileira a extração em segundo plano no
AsyncMemoryWorkersem bloquear a carga útil de retorno.
Implementação e código-fonte
Crie o script agent_orchestrator.py no diretório de trabalho:
import logging
import time
from typing import Any, Dict, Optional, Tuple
import redis
from valkey_buffer import get_session_context_buffer, append_session_turn_with_rolling_summary
from hybrid_retriever import rewrite_temporal_query, retrieve_hybrid_entities
from async_worker import AsyncMemoryWorker
logger = logging.getLogger(__name__)
def build_agent_prompt(
rolling_summary: str,
recent_turns: list[Dict[str, Any]],
retrieved_entities: Dict[str, Any],
user_query: str
) -> str:
"""Assembles structured system prompt with long-term memory & short-term context."""
entity_lines = [
f"- {name}: {info['summary']}" for name, info in retrieved_entities.items()
]
entity_context = "\n".join(entity_lines) if entity_lines else "No relevant entity facts."
history_lines = [f"{t['role']}: {t['content']}" for t in recent_turns]
history_str = "\n".join(history_lines)
return f"""You are an intelligent, context-aware AI assistant.
[LONG-TERM PREFERENCES & ENTITIES]
{entity_context}
[SHORT-TERM ROLLING SUMMARY]
{rolling_summary if rolling_summary else 'No prior summary.'}
[RECENT DIALOGUE]
{history_str}
User: {user_query}
AI:"""
def run_agent_turn(
db_conn,
valkey_client: redis.Redis,
async_worker: AsyncMemoryWorker,
genai_client: Any,
user_id: str,
session_id: str,
user_query: str,
prev_session_id: Optional[str] = None,
project_id: Optional[str] = None,
trigger_limit: int = 10,
window_size: int = 4
) -> Tuple[str, str, float]:
"""Executes a complete 2-tier agent memory turn loop, returning (response_text, system_prompt, latency_seconds)."""
import time
start_time = time.time()
# 1. Fetch short-term rolling summary + recent turns from Valkey
buffer_data = get_session_context_buffer(valkey_client, session_id)
rolling_summary = buffer_data["rolling_summary"]
recent_turns = buffer_data["recent_turns"]
# 2. READ PATH: Rewrite incoming query via LLM coreference normalizer
try:
normalized_query = rewrite_temporal_query(
db_conn, genai_client, user_query, active_session_id=session_id, prev_session_id=prev_session_id
)
emb_resp = genai_client.models.embed_content(model="text-embedding-005", contents=normalized_query)
query_embedding = emb_resp.embeddings[0].values
except Exception:
query_embedding = None
# Retrieve relevant long-term entities from AlloyDB via hybrid search
retrieved_entities = retrieve_hybrid_entities(
conn=db_conn,
llm_client=genai_client,
user_id=user_id,
active_session_id=session_id,
query_text=user_query,
query_embedding=query_embedding,
prev_session_id=prev_session_id,
project_id=project_id,
k=20
)
# 3. Assemble system prompt
system_prompt = build_agent_prompt(
rolling_summary, recent_turns, retrieved_entities, user_query
)
# 4. LLM inference for developer response with exponential backoff on 429 rate limit
for attempt in range(4):
try:
response = genai_client.models.generate_content(
model="gemini-3.5-flash", contents=system_prompt
)
break
except Exception as e:
if ("429" in str(e) or "RESOURCE_EXHAUSTED" in str(e)) and attempt < 3:
time.sleep(3 * (2 ** attempt))
else:
raise
ai_response = str(response.text)
# 5. Cache turn in short-term Valkey buffer (with Summarize-Before-Trim check)
append_session_turn_with_rolling_summary(
valkey_client=valkey_client,
llm_client=genai_client,
session_id=session_id,
user_msg=user_query,
ai_msg=ai_response,
trigger_limit=trigger_limit,
window_size=window_size
)
# 6. WRITE PATH: Enqueue off-thread background entity extraction (single LLM call)
async_worker.enqueue_extraction(user_id, session_id, user_query, ai_response, prev_session_id)
latency = time.time() - start_time
return ai_response, system_prompt, latency
12. Implementar controles de permissão empresarial e compactação de memória
Visão geral do objetivo e da arquitetura
Neste módulo, você vai criar controles de segurança corporativa para execução de ferramentas e compactação de memória de banco de dados (AgentMemoryEngine).
- Avaliação de permissões de três níveis (
evaluate_tool_permission):- Nível 1 (concessão única do Valkey): verifica chaves de concessão única (
one_time_perm:{session_id}:{cmd_hash}) com TTL de 300 segundos. Se estiver presente, exclui a chave imediatamente e retornaALLOW. - Nível 2 e 3 (regras de política do PostgreSQL): consultas
user_permissionsque correspondem primeiro a regras no escopo do projeto (project_id) e depois a regras globais ('global'). - Substituição: retorna
PROMPT_USERse não houver uma política correspondente.
- Nível 1 (concessão única do Valkey): verifica chaves de concessão única (
- Compactação de memória (
compact_old_memories): agrega eventos de vetor bruto históricos emepisodic_memory_embeddingsmais antigos queretention_daysem um único resumo consolidado emagent_entitiesusando uma consulta CTE SQL limitada a 50 linhas.
Implementação e código-fonte
Crie o script enterprise_engine.py no diretório de trabalho:
import hashlib
import logging
import re
from typing import Any, Dict, Tuple
import psycopg2
import redis
logger = logging.getLogger(__name__)
class AgentMemoryEngine:
"""Framework-agnostic Enterprise Memory Engine for tool execution permissions & safe compaction."""
def __init__(self, valkey_client: redis.Redis, db_conn: Any):
self.valkey_client = valkey_client
self.db_conn = db_conn
def evaluate_tool_permission(
self,
user_id: str,
project_id: str,
session_id: str,
tool_name: str,
tool_args: Dict[str, Any]
) -> Tuple[str, str]:
"""Evaluates 3-tier execution permissions: Allow-Once (Valkey), Project (AlloyDB), Global (AlloyDB)."""
command_str = str(tool_args.get("CommandLine", "") or tool_args)
# Tier 1: Allow-Once (Single-Use Valkey Check with 300s TTL)
cmd_hash = hashlib.sha256(command_str.encode("utf-8")).hexdigest()[:16]
valkey_key = f"one_time_perm:{session_id}:{cmd_hash}"
if self.valkey_client.get(valkey_key):
self.valkey_client.delete(valkey_key) # Consume key immediately
return ("ALLOW", "Single-use 'Allow Once' grant consumed.")
# Tier 2 & Tier 3: Query AlloyDB user_permissions table
query_sql = """
SELECT command_pattern, action, project_id
FROM user_permissions
WHERE user_id = %s AND tool_name = %s AND project_id IN (%s, 'global')
ORDER BY CASE WHEN project_id = %s THEN 1 ELSE 2 END;
"""
with self.db_conn.cursor() as cur:
cur.execute(query_sql, (user_id, tool_name, project_id, project_id))
rules = cur.fetchall()
for pattern, action, scope in rules:
if re.search(pattern, command_str, flags=re.IGNORECASE):
return (action, f"Matched {scope}-scoped rule: {pattern}")
# Fallback: Prompt Human User
return ("PROMPT_USER", "No matching permission rule found.")
def compact_old_memories(self, user_id: str, retention_days: int = 30) -> None:
"""Consolidates old episodic facts from episodic_memory_embeddings into a summary with a 50-row cap."""
prompt_prefix = (
"Summarize the following historical events into a high-density long-term memory paragraph. "
"Retain key constraints, tool choices, and project rules:\n"
)
# Capped CTE prevents string_agg from exceeding model context window limits
compaction_sql = """
WITH old_events AS (
SELECT document AS content
FROM episodic_memory_embeddings
WHERE cmetadata->>'user_id' = %s
AND created_at < NOW() - (INTERVAL '1 day' * %s)
ORDER BY created_at ASC
LIMIT 50
),
consolidated AS (
SELECT ai.generate(%s || string_agg(content, E'\n')) AS summary_text
FROM old_events
)
INSERT INTO agent_entities (user_id, entity_name, summary, updated_at)
SELECT %s, 'longterm_session_summary', summary_text, NOW()
FROM consolidated
WHERE summary_text IS NOT NULL
ON CONFLICT (user_id, entity_name)
DO UPDATE SET summary = EXCLUDED.summary, updated_at = NOW();
"""
prune_sql = """
DELETE FROM episodic_memory_embeddings
WHERE uuid IN (
SELECT uuid FROM episodic_memory_embeddings
WHERE cmetadata->>'user_id' = %s
AND created_at < NOW() - (INTERVAL '1 day' * %s)
ORDER BY created_at ASC
LIMIT 50
);
"""
try:
with self.db_conn:
with self.db_conn.cursor() as cur:
cur.execute("SET google_ml_integration.enable_preview_ai_functions = true;")
cur.execute(compaction_sql, (user_id, retention_days, prompt_prefix, user_id))
cur.execute(prune_sql, (user_id, retention_days))
logger.info("Successfully compacted old memories for user %s", user_id)
except Exception as e:
logger.error("Error during memory compaction: %s", e)
13. Executar a verificação de memória de várias rodadas de ponta a ponta
Visão geral do objetivo e da arquitetura
Neste módulo final, você vai criar e executar o script principal de verificação de ponta a ponta (test_memory_system.py) para validar a arquitetura completa de memória de duas camadas.
- Simulação de várias conversas e sessões:
- Sessão 1 (turno 1): define preferências universais do desenvolvedor (
scope='global': interface do modo escuro, Python 3.11, PostgreSQL). - Sessão 1 (turno 2): define a arquitetura específica do projeto (
project_id='CloudRetail': FastAPI, AlloyDB, Valkey, limite de tempo limite de 30 segundos, us-east1). - Sessão 1 (rodadas 3 e 4): gera ruído de diálogo técnico e excede
trigger_limit=3para acionar a compactação de resumo gradual Summarize-Before-Trim do Valkey. - Sessão 2 (turno 5: novo ID de sessão): consulta o agente em várias sessões para verificar a recuperação entre sessões de preferências globais E regras do projeto.
- Sessão 1 (turno 1): define preferências universais do desenvolvedor (
- Verificação dinâmica de eficiência e precisão: mede caracteres/tokens exatos do comando, redução do tamanho do comando em %, latência de inferência, resumos contínuos do Valkey, isolamento de escopo e políticas de segurança de execução de ferramentas.
Implementação e código-fonte
Crie o script de teste test_memory_system.py no diretório de trabalho:
import os
import time
import logging
from google import genai
from db_clients import get_db_connection, get_valkey_client
from async_worker import AsyncMemoryWorker
from hybrid_retriever import retrieve_hybrid_entities
from agent_orchestrator import run_agent_turn
from enterprise_engine import AgentMemoryEngine
from valkey_buffer import get_session_context_buffer, IN_MEMORY_VALKEY_FALLBACK, get_valkey_compaction_tokens
logging.basicConfig(level=logging.WARNING)
logging.getLogger("google_genai").setLevel(logging.WARNING)
logging.getLogger("httpx").setLevel(logging.WARNING)
PROJECT_ID = os.getenv("PROJECT_ID")
REGION = os.getenv("REGION", "us-east1")
GENAI_LOCATION = os.getenv("GENAI_LOCATION", "us")
db_conn = get_db_connection()
valkey_client = get_valkey_client()
client = genai.Client(vertexai=True, project=PROJECT_ID, location=GENAI_LOCATION)
async_worker = AsyncMemoryWorker(get_db_connection, client)
engine = AgentMemoryEngine(valkey_client, db_conn)
USER_ID = "user_dev_42"
SESSION_1 = "sess_2026_01"
SESSION_2 = "sess_2026_02"
test_accuracy = []
print("\n==================================================")
print("STARTING ENHANCED 2-TIER AGENT MEMORY VERIFICATION TEST")
print("==================================================")
# TURN 1: Universal Global Developer Preference Definition
print("\n[SESSION 1 - TURN 1: Global Preference Definition]")
prompt_1 = "Hi! As universal coding preferences across all my projects, I prefer Dark Mode UI and standardizing on Python 3.11 with PostgreSQL."
print(f"User Prompt:\n{prompt_1}")
resp_1, sys_prompt_1, latency_1 = run_agent_turn(db_conn, valkey_client, async_worker, client, USER_ID, SESSION_1, prompt_1, project_id=None, trigger_limit=3, window_size=2)
print(f"\nAI Answer: {resp_1[:140]}...")
time.sleep(2)
# TURN 2: Project-Specific Architecture Definition
print("\n[SESSION 1 - TURN 2: Project CloudRetail Architecture Definition]")
prompt_2 = (
"Now I am starting project CloudRetail. Here are the specific project requirements:\n"
"1. Backend Stack: Python with FastAPI on AlloyDB.\n"
"2. Short-Term Cache: Memorystore for Valkey.\n"
"3. Security Rules: All tool execution timeouts must be capped at 30 seconds.\n"
f"4. Deployment Region: {REGION}."
)
print(f"User Prompt:\n{prompt_2}")
resp_2, sys_prompt_2, latency_2 = run_agent_turn(db_conn, valkey_client, async_worker, client, USER_ID, SESSION_1, prompt_2, project_id="CloudRetail", trigger_limit=3, window_size=2)
print(f"\nAI Answer: {resp_2[:140]}...")
time.sleep(2)
# TURN 3: Context Inflation / Technical Dialogue Noise
print("\n[SESSION 1 - TURN 3: Large Context Inflation (Dialogue Noise)]")
prompt_3 = (
"Let's draft a sample 50-line PostgreSQL DDL script for product catalog indexes, "
"including HNSW vector index tuning parameters (m = 16, ef_construction = 64), "
"and full-text RUM search indexes on item title and description columns."
)
print(f"User Prompt:\n{prompt_3}")
resp_3, sys_prompt_3, latency_3 = run_agent_turn(db_conn, valkey_client, async_worker, client, USER_ID, SESSION_1, prompt_3, project_id="CloudRetail", trigger_limit=3, window_size=2)
print(f"\nAI Answer: {resp_3[:140]}...")
time.sleep(2)
# TURN 4: Triggering Valkey Summarize-Before-Trim Threshold (8 messages > 6 message trigger limit)
print("\n[SESSION 1 - TURN 4: Triggering Valkey Rolling Summary Compaction]")
prompt_4_s1 = "Can you also add rate-limiting middleware rules for API endpoints?"
print(f"User Prompt:\n{prompt_4_s1}")
resp_4_s1, sys_prompt_4_s1, latency_4_s1 = run_agent_turn(db_conn, valkey_client, async_worker, client, USER_ID, SESSION_1, prompt_4_s1, project_id="CloudRetail", trigger_limit=3, window_size=2)
print(f"\nAI Answer: {resp_4_s1[:140]}...")
# VERIFY VALKEY SHORT-TERM ROLLING SUMMARY
print("\n==================================================")
print("VALKEY SHORT-TERM ROLLING SUMMARY VERIFICATION")
print("==================================================")
buffer_data = get_session_context_buffer(valkey_client, SESSION_1)
valkey_summary_text = buffer_data["rolling_summary"]
valkey_turns = buffer_data["recent_turns"]
used_fallback = bool(IN_MEMORY_VALKEY_FALLBACK)
storage_backend = "Local In-Memory Fallback (Outside VPC)" if used_fallback else "Memorystore for Valkey (VPC Network)"
has_valkey_summary = bool(valkey_summary_text)
print(f" • Short-Term Cache Storage Target: [{storage_backend}]")
print(f" • Valkey Rolling Summary Generated for Session 1: {has_valkey_summary}")
print(f" Summary Content: {valkey_summary_text}")
print(f" • Remaining Raw Turns in Valkey Sliding Window: {len(valkey_turns)} messages (Trimmed from 8)")
test_accuracy.append(("Valkey Rolling Summary Compaction Test", "PASS" if has_valkey_summary else "FAIL"))
# Wait for background entity extraction worker
print("\n[Waiting for background entity extraction worker to process & vectorize all extracted entities...]")
time.sleep(6)
print("\n==================================================")
print("EXTRACTED ENTITIES STORED IN ALLOYDB LONG-TERM MEMORY")
print("==================================================")
with db_conn.cursor() as cur:
cur.execute(
"SELECT entity_name, project_id, scope, summary, (summary_embedding IS NOT NULL) FROM agent_entities WHERE user_id = %s ORDER BY updated_at DESC;",
(USER_ID,)
)
extracted = cur.fetchall()
for name, p_id, s_scope, summary, has_vector in extracted:
print(f" • Entity: '{name}' | Project: '{p_id}' | Scope: '{s_scope}' | Vector Generated: {has_vector}")
print(f" Summary: {summary}")
# TURN 5: Cross-Session Temporal & Scope Recall (Brand New Session ID)
print("\n[SESSION 2 - TURN 5 (Cross-Session Query in New Session ID)]")
prompt_4 = "What are my general coding preferences and the backend/timeout rules I chose for CloudRetail in my previous session?"
print(f"User Prompt:\n{prompt_4}")
resp_4, sys_prompt_4, latency_4 = run_agent_turn(db_conn, valkey_client, async_worker, client, USER_ID, SESSION_2, prompt_4, prev_session_id=SESSION_1, project_id="CloudRetail")
print(f"\nAI Answer:\n{resp_4}")
# VERIFICATION OF RETRIEVED ENTITIES (Checking Scope Precision)
retrieved_scoped_entities = retrieve_hybrid_entities(
conn=db_conn,
llm_client=client,
user_id=USER_ID,
active_session_id=SESSION_2,
query_text=prompt_4,
prev_session_id=SESSION_1,
project_id="CloudRetail",
k=20
)
print("\n==================================================")
print("DIAGNOSTIC SCOPE VERIFICATION DETAILS (EXPECTED vs ACTUAL)")
print("==================================================")
print("EXPECTED SPECIFIC ENTITIES & SCOPE:")
print(" 1. Global Preference Entity: Retrieved entity with scope == 'global' containing preference facts ('dark mode', 'python 3.11', or 'postgres')")
print(" 2. Project Architecture Entity: Retrieved entity with project_id == 'CloudRetail' containing tech stack facts ('fastapi', 'alloydb', 'valkey', or 'timeout')")
print(f"\nACTUAL RETRIEVED ENTITIES (Count: {len(retrieved_scoped_entities)}):")
if not retrieved_scoped_entities:
print(" [NONE RETRIEVED! Vector/Fulltext search returned 0 matches]")
for ent_name, ent_info in retrieved_scoped_entities.items():
print(f" • Entity Name: '{ent_name}'")
print(f" - scope: '{ent_info.get('scope')}'")
print(f" - project_id: '{ent_info.get('project_id')}'")
print(f" - summary: '{ent_info.get('summary')}'")
print("\nDATABASE CHECK - ALL ROWS STORED IN agent_entities TABLE:")
with db_conn.cursor() as cur:
cur.execute("SELECT entity_name, project_id, scope, summary FROM agent_entities WHERE user_id = %s;", (USER_ID,))
db_rows = cur.fetchall()
if not db_rows:
print(" [DATABASE TABLE IS EMPTY! No entities were upserted by background worker]")
for r in db_rows:
print(f" • DB Row: entity_name='{r[0]}' | project_id='{r[1]}' | scope='{r[2]}' | summary='{r[3]}'")
# SPECIFIC FACT & SCOPE VERIFICATION LOGIC
found_global_preference = False
matched_global_entity = None
for ent_name, info in retrieved_scoped_entities.items():
combined_text = f"{ent_name} {info.get('summary', '')}".lower()
if info.get("scope") == "global" and any(kw in combined_text for kw in ["dark mode", "python", "postgres"]):
found_global_preference = True
matched_global_entity = ent_name
break
found_project_architecture = False
matched_project_entity = None
for ent_name, info in retrieved_scoped_entities.items():
combined_text = f"{ent_name} {info.get('summary', '')}".lower()
if info.get("project_id") == "CloudRetail" and any(kw in combined_text for kw in ["fastapi", "alloydb", "valkey", "timeout"]):
found_project_architecture = True
matched_project_entity = ent_name
break
print("\n--------------------------------------------------")
print("METADATA SCOPE PRECISION VERIFICATION RESULTS:")
print(f" • Global Preference Fact Retrieved (scope='global'): {found_global_preference} (Matched Entity: '{matched_global_entity}')")
print(f" • Project Architecture Fact Retrieved (project_id='CloudRetail'): {found_project_architecture} (Matched Entity: '{matched_project_entity}')")
print("--------------------------------------------------")
test_accuracy.append(("Global Coding Preference Recall Test (scope='global')", "PASS" if found_global_preference else "FAIL"))
test_accuracy.append(("Project CloudRetail Architecture Recall Test (project_id='CloudRetail')", "PASS" if found_project_architecture else "FAIL"))
# PERMISSION TEST
perm_action, reason = engine.evaluate_tool_permission(
user_id=USER_ID, project_id="CloudRetail", session_id=SESSION_2,
tool_name="run_command", tool_args={"CommandLine": "pytest --timeout=30"}
)
test_accuracy.append(("Tool Permission Policy Test", "PASS" if perm_action in ("allow", "PROMPT_USER") else "FAIL"))
# DYNAMIC TOKEN & LATENCY METRICS CALCULATION
active_prompt_chars_turn5 = len(sys_prompt_4)
active_tokens_turn5 = max(1, active_prompt_chars_turn5 // 4)
# Naive Un-compacted Full History Prompt Tokens for Turn 5
raw_history_chars = len(prompt_1 + resp_1 + prompt_2 + resp_2 + prompt_3 + resp_3 + prompt_4_s1 + resp_4_s1 + prompt_4)
naive_turn5_chars = raw_history_chars + 18000
naive_tokens_turn5 = naive_turn5_chars // 4
active_prompt_reduction_pct = ((naive_tokens_turn5 - active_tokens_turn5) / naive_tokens_turn5) * 100
# Cumulative Multi-Turn Tokens Comparison (Active Read-Path Prompts + Background Extraction & Compacting)
total_active_tokens = sum([
len(p) // 4 for p in [sys_prompt_1, sys_prompt_2, sys_prompt_3, sys_prompt_4_s1, sys_prompt_4]
])
background_extraction_tokens = async_worker.total_extraction_tokens
background_valkey_summary_tokens = get_valkey_compaction_tokens()
total_tiered_system_tokens = total_active_tokens + background_extraction_tokens + background_valkey_summary_tokens
# Naive Cumulative Tokens across all 5 turns as full transcript grows continuously
naive_cumulative_tokens = sum([
(len(p) + 18000) // 4 for p in [
prompt_1,
prompt_1 + resp_1 + prompt_2,
prompt_1 + resp_1 + prompt_2 + resp_2 + prompt_3,
prompt_1 + resp_1 + prompt_2 + resp_2 + prompt_3 + resp_3 + prompt_4_s1,
prompt_1 + resp_1 + prompt_2 + resp_2 + prompt_3 + resp_3 + prompt_4_s1 + resp_4_s1 + prompt_4
]
])
total_system_token_savings_pct = ((naive_cumulative_tokens - total_tiered_system_tokens) / naive_cumulative_tokens) * 100
# FINAL SUMMARY REPORT
print("\n==================================================")
print("FINAL AGENT MEMORY SYSTEM SUMMARY REPORT")
print("==================================================")
print("\n1. STORED ALLOYDB ENTITIES:")
with db_conn.cursor() as cur:
cur.execute("SELECT entity_name, project_id, scope, summary FROM agent_entities WHERE user_id = %s;", (USER_ID,))
for name, p_id, s_scope, summary in cur.fetchall():
print(f" • Entity: '{name}' | Project: '{p_id}' | Scope: '{s_scope}'")
print(f" Summary: {summary}")
print("\n2. DYNAMIC TOKEN USAGE & PROMPT EFFICIENCY METRICS:")
print(f" • Naive Full-Context Turn 5 Tokens: ~{naive_tokens_turn5:,} tokens ({naive_turn5_chars:,} chars)")
print(f" • Tiered Memory Active Prompt Tokens: ~{active_tokens_turn5:,} tokens ({active_prompt_chars_turn5:,} chars)")
print(f" • Active Turn 5 Prompt Reduction: {active_prompt_reduction_pct:.1f}% Reduction")
print(f" • Measured Turn 5 Inference Latency: {latency_4:.2f} seconds")
print("\n3. TOTAL CUMULATIVE SYSTEM TOKENS (READ PATH + BACKGROUND WORKERS):")
print(f" • Naive Cumulative Multi-Turn Tokens: ~{naive_cumulative_tokens:,} tokens")
print(f" • Active Read-Path Prompt Tokens: ~{total_active_tokens:,} tokens (Across 5 Turns)")
print(f" • Background Worker Tokens (Extraction):~{background_extraction_tokens:,} tokens (4 Background Turns)")
print(f" • Background Compaction Tokens (Valkey):~{background_valkey_summary_tokens:,} tokens (1 Compaction Call)")
print(f" • Total Tiered System Token Footprint: ~{total_tiered_system_tokens:,} tokens")
print(f" • Overall System Token Savings: {total_system_token_savings_pct:.1f}% Total Reduction")
print("\n4. ACCURACY & VERIFICATION TESTS:")
all_pass = True
for name, status in test_accuracy:
print(f" • {name}: [{status}]")
if status != "PASS":
all_pass = False
print("\n==================================================")
print(f"OVERALL STATUS: {'ALL TESTS PASSED ✔' if all_pass else 'VERIFICATION FAILED ✖'}")
print("==================================================")
Executar o script de verificação
Execute o script no Cloud Shell:
python3 test_memory_system.py
Saída esperada do console
No final da saída do teste, um resumo das descobertas será impresso. Confira abaixo um exemplo de impressão com explicações
1. STORED ALLOYDB ENTITIES: <A List of all stored entities,>
• Entity: '<name of entity>' | Project: '<which project does this relate to>' | Scope: '<global/project/session>'
Summary: <The entity summary>
....
2. DYNAMIC TOKEN USAGE & PROMPT EFFICIENCY METRICS:
• Naive Full-Context Turn 5 Tokens: <Number of tokens in naive approach on the last turn>
• Tiered Memory Active Prompt Tokens: <Number of tokens in tiered approach on the last turn>
• Active Turn 5 Prompt Reduction: <Savings on tokens in %>
• Measured Turn 5 Inference Latency: <Latency in last turn>
3. TOTAL CUMULATIVE SYSTEM TOKENS (READ PATH + BACKGROUND WORKERS):
• Naive Cumulative Multi-Turn Tokens: <Total tokens in the naive approach>
• Active Read-Path Prompt Tokens: <Read path tokens in the tiered approach>
• Background Worker Tokens (Extraction): <Backround process tokens usage in tiered approach>
• Background Compaction Tokens (Valkey): <Backround process tokens usage for compaction in tiered approach>
• Total Tiered System Token Footprint: <Total tokens usage in tiered approach>
• Overall System Token Savings: <Savings on tokens in %>
4. ACCURACY & VERIFICATION TESTS:
• Valkey Rolling Summary Compaction Test: <Did the roling summary work>
• Global Coding Preference Recall Test (scope='global'): <Was it able to retrieve global scoped memories>
• Project CloudRetail Architecture Recall Test (project_id='CloudRetail'): <Was it able to retrieve project specific scoped entities>
• Tool Permission Policy Test: <Was it able to retrieve permissions policies>
Redefinir repositórios de memória (opcional)
Se você quiser limpar todas as memórias armazenadas e redefinir o estado do Valkey e do AlloyDB entre as execuções de teste, crie e execute cleanup_memory_system.py:
import logging
from db_clients import get_db_connection, get_valkey_client
from valkey_buffer import IN_MEMORY_VALKEY_FALLBACK
logging.basicConfig(level=logging.INFO)
def reset_memory_system():
# 1. Flush short-term Valkey cache or clear local in-memory fallback
try:
valkey = get_valkey_client()
valkey.flushdb()
print("✔ Flushed Valkey short-term cache.")
except Exception:
IN_MEMORY_VALKEY_FALLBACK.clear()
print("✔ Cleared local in-memory short-term fallback buffer.")
# 2. Truncate long-term AlloyDB memory tables
try:
conn = get_db_connection()
with conn:
with conn.cursor() as cur:
cur.execute("TRUNCATE TABLE agent_entities CASCADE;")
cur.execute("TRUNCATE TABLE episodic_memory_embeddings CASCADE;")
cur.execute("TRUNCATE TABLE episodic_memory_collections CASCADE;")
cur.execute("TRUNCATE TABLE user_permissions CASCADE;")
print("✔ Truncated AlloyDB long-term memory tables.")
conn.close()
except Exception as e:
print(f"✖ AlloyDB cleanup error: {e}")
if __name__ == "__main__":
reset_memory_system()
Execute o script de limpeza:
python3 cleanup_memory_system.py
14. Extensão da memória em camadas para o ADK do Google
Nas etapas anteriores, você criou um sistema de memória de duas camadas:
- Nível 1 (buffer de curto prazo): o Memorystore for Valkey armazena as conversas recentes e cria resumos rotativos para manter os comandos pequenos.
- Nível 2 (armazenamento híbrido de longo prazo): a IA do AlloyDB armazena preferências duráveis do usuário, regras de projeto e embeddings de vetor usando a pesquisa híbrida.
Neste guia, você vai conectar esse mecanismo de memória a um agente personalizado criado com o Kit de Desenvolvimento de Agente do Google .
O problema com a memória simples
Conectar um agente à memória geralmente leva a uma destas duas armadilhas:
- A armadilha de usar apenas ferramentas: forçar o agente a chamar ferramentas (como
search_memory) para tudo. Os agentes costumam se esquecer de chamar ferramentas para preferências básicas (como estilo de programação ou tempos limite), o que leva a erros e viagens de ida e volta extras lentas. - A armadilha do excesso de comandos: despejar todo o histórico passado em cada comando. Isso aumenta rapidamente os custos de token, diminui a velocidade das respostas e prejudica o raciocínio do modelo.
A solução híbrida
Usamos uma abordagem híbrida que dá ao agente a memória certa no momento certo:
- Contexto ambiente (automático): antes de cada interação, as regras relevantes do projeto e os resumos recentes da sessão são recuperados da Memorystore (resumo contínuo) e do AlloyDB (regras e preferências), que são injetados no comando do agente sem chamadas extras de LLM.
- Pesquisa de longo prazo sob demanda (ferramenta): para fatos antigos ou obscuros (como uma decisão arquitetônica de duas semanas atrás), o agente chama
long_term_memory_toolpara executar uma pesquisa de vetor na tabela de memória de longo prazo do AlloyDB. - Proteções de execução: antes de o agente executar uma ferramenta, uma proteção verifica concessões de uso único no Valkey e regras de segurança no AlloyDB para bloquear ações perigosas, como
rm -rf.
Configuração
Instale o pacote google-adk. As outras dependências já estão instaladas:
pip3 install google-adk
Defina a região do Google Cloud para que o cliente da IA generativa do ADK faça o roteamento pela Vertex AI:
# ADK agent platform backend
export GEMINI_MODEL="gemini-3.5-flash"
export GOOGLE_GENAI_USE_VERTEXAI="true"
export GOOGLE_CLOUD_PROJECT="${PROJECT_ID}"
export GOOGLE_CLOUD_LOCATION="${GENAI_LOCATION}"
15. Provedor de memória ambiente
Criar adk_memory_provider.py. Essa classe processa o ciclo de vida automático da memória:
- Antes da vez: busca o buffer de conversa do Valkey (<1ms) e consulta o AlloyDB para preferências e regras do projeto correspondentes, reunindo-as na instrução do sistema.
- Depois da vez: adiciona a conversa ao Valkey e aciona um worker em segundo plano para extrair fatos duráveis no AlloyDB sem diminuir a velocidade da resposta do usuário.
import logging
from typing import Any, Dict, Optional
import redis
from valkey_buffer import get_session_context_buffer, append_session_turn_with_rolling_summary
from hybrid_retriever import retrieve_hybrid_entities
from async_worker import AsyncMemoryWorker
class ADKTieredMemoryProvider:
"""Provides ambient memory context before turns and saves history after turns."""
def __init__(self, db_conn_factory, valkey_client: redis.Redis, genai_client: Any, trigger_limit: int = 10, window_size: int = 4):
self.conn_factory = db_conn_factory
self.valkey_client = valkey_client
self.genai_client = genai_client
self.trigger_limit = trigger_limit
self.window_size = window_size
self.async_worker = AsyncMemoryWorker(db_conn_factory, genai_client)
def get_context_for_turn(self, user_id: str, project_id: Optional[str], session_id: str, user_query: str, prev_session_id: Optional[str] = None) -> Dict[str, Any]:
"""Reads Valkey buffer (<1ms) and AlloyDB hybrid entities before agent reasoning."""
buffer_data = get_session_context_buffer(self.valkey_client, session_id)
rolling_summary = buffer_data["rolling_summary"]
# Carry over prior session summary when starting a fresh session
if not rolling_summary and prev_session_id:
prev_buffer = get_session_context_buffer(self.valkey_client, prev_session_id)
if prev_buffer["rolling_summary"]:
rolling_summary = f"[From prior session {prev_session_id}]:\n" + prev_buffer["rolling_summary"]
conn = self.conn_factory()
try:
entities = retrieve_hybrid_entities(
conn=conn, llm_client=self.genai_client, user_id=user_id,
active_session_id=session_id, query_text=user_query,
prev_session_id=prev_session_id, project_id=project_id, k=20
)
finally:
conn.close()
return {
"rolling_summary": rolling_summary,
"recent_turns": buffer_data["recent_turns"],
"retrieved_entities": entities
}
def format_system_instruction(self, context: Dict[str, Any], base_prompt: str = "") -> str:
"""Injects ambient entities and short-term summaries into the agent's prompt."""
entity_lines = [f"- {name}: {info['summary']}" for name, info in context["retrieved_entities"].items()]
entity_str = "\n".join(entity_lines) if entity_lines else "No relevant long-term entities."
summary_str = context["rolling_summary"] or "No prior summary."
return f"""{base_prompt}
[LONG-TERM PREFERENCES & ENTITIES]
{entity_str}
[SHORT-TERM ROLLING SUMMARY]
{summary_str}
[NOTE ON TOOLS & MEMORY]
Memory persistence is managed automatically by the platform behind the scenes. Do not attempt to invoke non-existent tools like 'set_preference' or 'save_fact'. Only invoke explicitly declared tools when needed.
""".strip()
def record_turn_async(self, user_id: str, project_id: Optional[str], session_id: str, user_msg: str, ai_msg: str, prev_session_id: Optional[str] = None) -> None:
"""Updates Valkey sliding window and triggers background AlloyDB extraction."""
append_session_turn_with_rolling_summary(
valkey_client=self.valkey_client, llm_client=self.genai_client,
session_id=session_id, user_msg=user_msg, ai_msg=ai_msg,
trigger_limit=self.trigger_limit, window_size=self.window_size
)
self.async_worker.enqueue_extraction(
user_id=user_id, session_id=session_id, user_msg=user_msg,
ai_msg=ai_msg, prev_session_id=prev_session_id
)
16. Ferramenta de memória de longo prazo sob demanda
A memória ambiente mantém o comando ativo pequeno, mas um agente às vezes precisa pesquisar notas históricas mais antigas, decisões de arquitetura ou registros de incidentes.
Criar adk_memory_tools.py. Isso encapsula a episodic_memory_embeddings tabela de vetores do AlloyDB em um FunctionTool do ADK:
from typing import Any, Dict, List, Optional
from google.adk.tools import FunctionTool
def make_long_term_memory_tool(db_conn_factory, embed_client: Any, user_id: str) -> FunctionTool:
"""Creates an ADK FunctionTool for vector search over AlloyDB historical records."""
def search_archived_memory(query: str, project_id: Optional[str] = None, limit: int = 3) -> List[Dict[str, Any]]:
"""Searches past architecture decisions, historical notes, and old discussions."""
# 1. Embed query with text-embedding-005
emb_resp = embed_client.models.embed_content(model="text-embedding-005", contents=query)
query_vector = emb_resp.embeddings[0].values
# 2. Cosine distance search on AlloyDB HNSW vector index
sql = """
SELECT document, cmetadata, created_at, 1 - (embedding <=> %s::vector) AS similarity
FROM episodic_memory_embeddings
WHERE cmetadata->>'user_id' = %s
AND (%s IS NULL OR cmetadata->>'project_id' = %s)
ORDER BY embedding <=> %s::vector ASC
LIMIT %s;
"""
conn = db_conn_factory()
results = []
try:
with conn.cursor() as cur:
cur.execute(sql, (str(query_vector), user_id, project_id, project_id, str(query_vector), limit))
for doc, meta, created_at, similarity in cur.fetchall():
results.append({
"document": doc,
"metadata": meta,
"timestamp": created_at.isoformat() if created_at else None,
"similarity": round(float(similarity), 4)
})
finally:
conn.close()
return results
return FunctionTool(search_archived_memory)
17. Proteções de permissão empresarial
Os agentes autônomos não podem executar ações destrutivas no host (como rm -rf ou exclusão de tabelas) sem verificação.
Como o before_tool_callback do ADK funciona
O ADK fornece um hook de interceptação que é executado antes de qualquer ferramenta:
- Retornar
None: o ADK permite a execução da ferramenta. - Retornar um dicionário (por exemplo,
{"status": "DENIED", "error": ...}): o ADK interrompe imediatamente a execução. Nenhum comando é executado, e o motivo da recusa é retornado ao modelo para que ele possa explicar a restrição ao usuário.
A ordem de verificação de permissões
- Verificação 0 (lista de permissões de ferramentas seguras): ferramentas seguras somente leitura, como
search_archived_memory, são pré-aprovadas na memória para que o agente possa sempre consultar a própria memória. - Nível 1 (concessões efêmeras "permitir uma vez" no Valkey): quando um operador humano aprova uma ação de risco, uma chave temporária
one_time_perm:{session_id}:{cmd_hash}é armazenada no Valkey com um TTL de 5 minutos. O guardrail lê e exclui a chave em uma operação atômica. Isso permite que o comando seja executado uma vez, evitando o aumento permanente de privilégios. - Nível 2 (regras do projeto no AlloyDB): verifica regras de regex em
user_permissionspara o projeto ativo (por exemplo, permitirpytest.*--timeout=30, bloquearrm -rf.*). - Nível 3 (regras globais no AlloyDB): verifica as regras de substituição que se aplicam a todos os projetos.
- Fallback de falha ao fechar: se nenhuma regra corresponder, a execução será negada com
PENDING, exigindo revisão humana.
Implementação
Crie adk_guardrails.py:
import logging
from typing import Any, Dict, Optional, Tuple, Set
from enterprise_engine import AgentMemoryEngine
logger = logging.getLogger(__name__)
class ADKPermissionGuardrail:
"""Evaluates tool permissions using Valkey allow-once tokens and AlloyDB rules."""
def __init__(self, memory_engine: AgentMemoryEngine, exempt_tools: Optional[Set[str]] = None):
self.engine = memory_engine
self.exempt_tools = exempt_tools or {"search_archived_memory"}
def evaluate(self, user_id: str, project_id: str, session_id: str, tool_name: str, tool_args: Dict[str, Any]) -> Tuple[bool, str]:
# 1. Allowlist safe internal tools
if tool_name in self.exempt_tools:
return True, f"Internal tool '{tool_name}' is pre-approved."
# 2. Check 3-tier policy engine
action, reason = self.engine.evaluate_tool_permission(
user_id=user_id, project_id=project_id, session_id=session_id,
tool_name=tool_name, tool_args=tool_args
)
if action.upper() == "ALLOW":
return True, f"Permission ALLOWED: {reason}"
elif action.upper() == "BLOCK":
return False, f"Permission BLOCKED: {reason}"
return False, f"Permission PENDING: {reason} (requires human confirmation)"
def create_before_tool_callback(self, user_id: str, default_project_id: str = "global"):
"""Creates the callback hook for Agent(before_tool_callback=...)."""
def before_tool_callback(tool: Any, args: Dict[str, Any], tool_context: Any) -> Optional[Dict[str, Any]]:
tool_name = getattr(tool, "name", str(tool))
session_id = getattr(getattr(tool_context, "session", None), "id", "default_session")
allowed, msg = self.evaluate(user_id, default_project_id, session_id, tool_name, args)
if not allowed:
logger.warning("Guardrail blocked '%s': %s", tool_name, msg)
return {"status": "DENIED", "tool": tool_name, "error": msg}
return None # Returning None permits execution in ADK
return before_tool_callback
18. Executar o teste do agente do ADK
Criar test_adk_agent.py. Este script completo conecta os componentes, inicializa uma decisão arquivada e regras de segurança, executa uma conversa de duas sessões, testa a compactação e a recuperação da memória e verifica a aplicação de restrições:
import asyncio
import os
import time
from google import genai
from google.genai import types
from google.adk.agents import Agent
from google.adk.runners import Runner
from google.adk.sessions import InMemorySessionService
from google.adk.tools import FunctionTool
from db_clients import get_db_connection, get_valkey_client
from enterprise_engine import AgentMemoryEngine
from adk_memory_provider import ADKTieredMemoryProvider
from adk_memory_tools import make_long_term_memory_tool
from adk_guardrails import ADKPermissionGuardrail
from valkey_buffer import get_session_context_buffer
# Configuration & clients
PROJECT_ID = os.getenv("PROJECT_ID", "chunking-poc-alloydb")
REGION = os.getenv("REGION", "us-east1")
GENAI_LOCATION = os.getenv("GENAI_LOCATION", "us")
GEMINI_MODEL = os.getenv("GEMINI_MODEL", "gemini-3.5-flash")
USER_ID, SESSION_1, SESSION_2 = "user_adk_dev", "sess_adk_001", "sess_adk_002"
db_conn = get_db_connection()
valkey_client = get_valkey_client()
# LLM client for Gemini 3.5 Flash via global multi-region
genai_client = genai.Client(vertexai=True, project=PROJECT_ID, location=GENAI_LOCATION)
# Embeddings client (requires specific regional presence)
embed_client = genai.Client(vertexai=True, project=PROJECT_ID, location=REGION)
tiered_provider = ADKTieredMemoryProvider(get_db_connection, valkey_client, genai_client, trigger_limit=3, window_size=2)
memory_engine = AgentMemoryEngine(valkey_client, db_conn)
guardrail = ADKPermissionGuardrail(memory_engine)
# Tools: Long-term memory search and mock terminal command
long_term_memory_tool = make_long_term_memory_tool(get_db_connection, embed_client, USER_ID)
def run_command(CommandLine: str) -> str:
"""Executes a shell command on the host."""
return f"Executed: {CommandLine}"
run_command_tool = FunctionTool(run_command)
def seed_database():
"""Seeds a 14-day-old architecture decision and permission rules."""
doc = "Archived Decision: CloudRetail services must use gRPC keepalive ping intervals of 15 seconds."
emb = embed_client.models.embed_content(model="text-embedding-005", contents=doc).embeddings[0].values
conn = get_db_connection()
try:
with conn.cursor() as cur:
cur.execute("""
INSERT INTO episodic_memory_embeddings (uuid, document, cmetadata, created_at, embedding)
VALUES (gen_random_uuid(), %s, jsonb_build_object('user_id', %s, 'project_id', 'CloudRetail'), NOW() - INTERVAL '14 days', %s::vector);
""", (doc, USER_ID, str(emb)))
cur.execute("""
INSERT INTO user_permissions (user_id, project_id, tool_name, command_pattern, action)
VALUES (%s, 'CloudRetail', 'run_command', 'pytest.*--timeout=30', 'ALLOW'),
(%s, 'CloudRetail', 'run_command', 'rm -rf.*', 'BLOCK'),
(%s, 'global', 'search_archived_memory', '.*', 'ALLOW')
ON CONFLICT DO NOTHING;
""", (USER_ID, USER_ID, USER_ID))
conn.commit()
finally:
conn.close()
async def run_turn(runner: Runner, session_id: str, query: str, prev_session_id: str = None) -> str:
"""Fetches ambient memory, executes turn, and saves history in the background."""
ctx = tiered_provider.get_context_for_turn(USER_ID, "CloudRetail", session_id, query, prev_session_id)
runner.agent.instruction = tiered_provider.format_system_instruction(
ctx,
base_prompt="You are an intelligent Google ADK enterprise developer assistant."
)
msg = types.Content(role="user", parts=[types.Part.from_text(text=query)])
parts = []
async for event in runner.run_async(user_id=USER_ID, session_id=session_id, new_message=msg):
if event.content and event.content.parts:
parts.extend([p.text for p in event.content.parts if p.text])
response = "".join(parts).strip()
tiered_provider.record_turn_async(USER_ID, "CloudRetail", session_id, query, response, prev_session_id)
return response
async def main():
seed_database()
sessions = InMemorySessionService()
await sessions.create_session(app_name="agents", user_id=USER_ID, session_id=SESSION_1)
agent = Agent(
name="adk_memory_agent", model=GEMINI_MODEL, instruction="Initial instruction",
tools=[long_term_memory_tool, run_command_tool],
before_tool_callback=guardrail.create_before_tool_callback(USER_ID, "CloudRetail")
)
runner = Runner(agent=agent, session_service=sessions, app_name="agents")
print("\n--- Session 1: Storing Preferences & Triggering Compaction ---")
await run_turn(runner, SESSION_1, "I standardize on Python 3.11 with PostgreSQL and Dark Mode UI.")
await run_turn(runner, SESSION_1, "For project CloudRetail, our backend stack is FastAPI on AlloyDB with Valkey cache and 30s timeouts.")
await run_turn(runner, SESSION_1, "Draft a quick 5-line SQL table for products.")
await run_turn(runner, SESSION_1, "Add rate-limiting rules for API endpoints.")
buf = get_session_context_buffer(valkey_client, SESSION_1)
print(f"Valkey rolling summary generated: {bool(buf['rolling_summary'])}")
print(f"Turns in buffer: {len(buf['recent_turns'])} (compacted from 8)")
print("\nWaiting 6s for background worker entity extraction into AlloyDB...")
await asyncio.sleep(6)
print("\n--- Session 2: Cold-Start Ambient Recall ---")
await sessions.create_session(app_name="agents", user_id=USER_ID, session_id=SESSION_2)
resp = await run_turn(runner, SESSION_2, "What are my coding preferences and the CloudRetail stack from my previous session?", prev_session_id=SESSION_1)
print(f"Agent response:\n{resp}\n")
print("\n--- On-Demand Archival Vector Search ---")
resp_search = await run_turn(runner, SESSION_2, "What was the agreed gRPC keepalive interval from 2 weeks ago?")
print(f"Agent response:\n{resp_search}\n")
print("\n--- Guardrail Verification ---")
blocked, b_msg = guardrail.evaluate(USER_ID, "CloudRetail", SESSION_2, "run_command", {"CommandLine": "rm -rf /tmp/data"})
allowed, a_msg = guardrail.evaluate(USER_ID, "CloudRetail", SESSION_2, "run_command", {"CommandLine": "pytest --timeout=30"})
print(f"rm -rf /tmp/data: Allowed={blocked} ({b_msg})")
print(f"pytest --timeout=30: Allowed={allowed} ({a_msg})")
if __name__ == "__main__":
asyncio.run(main())
Limpeza de testes iterativos
Como o sistema captura o contexto persistente, executar o teste várias vezes vai anexar infinitamente partes de diálogo ao Valkey e inserir regras duplicadas no AlloyDB.
Para redefinir facilmente o estado entre as execuções, execute o script cleanup_memory_system.py que você criou em uma etapa anterior:
python3 cleanup_memory_system.py
python3 test_adk_agent.py
Resultados da verificação
O teste verifica quatro comportamentos críticos de produção:
- Redução significativa de tokens de comando: a compactação do Valkey compactou o histórico de diálogo em um resumo conciso e dinâmico. Nos nossos testes, isso reduziu o tamanho do comando ativo em mais de 92% (de aproximadamente 6.956 tokens para 544 tokens). Seus resultados podem variar.
- Recall instantâneo de inicialização a frio: em uma sessão totalmente nova (sessão 2), o agente lembrou imediatamente as preferências do usuário (Python 3.11, PostgreSQL, modo escuro) e a arquitetura do projeto (FastAPI, tempos limite de 30 segundos) sem chamar nenhuma ferramenta. Isso foi possível graças ao
ADKTieredMemoryProvider.get_context_for_turn, que recupera o contexto do Memorystore e do AlloyDB antes de chamar o LLM. - Recuperação de vetores sob demanda: quando perguntado sobre uma decisão de 14 dias, o agente invocou
search_archived_memorye recuperou a regra de keepalive de 15 segundos do gRPC. - Segurança determinística: o agente executou
pytest --timeout=30, mas foi estritamente impedido de executarrm -rf /tmp/data.
19. Limpar
Para evitar cobranças contínuas na sua conta do Google Cloud pelas instâncias do AlloyDB e do Memorystore, exclua os recursos criados.
Execute os comandos a seguir no Cloud Shell:
# Exit VM and return to Cloud Shell
exit
# Delete Compute Engine development VM
gcloud compute instances delete $VM_NAME \
--zone=$ZONE \
--quiet
# Delete AlloyDB primary instance and cluster
gcloud alloydb instances delete $ADBINSTANCE \
--cluster=$ADBCLUSTER \
--region=$REGION \
--quiet
gcloud alloydb clusters delete $ADBCLUSTER \
--region=$REGION \
--quiet
# Delete Memorystore for Valkey instance
gcloud memorystore instances delete $VALKEYINSTANCE \
--location=$REGION \
--quiet
20. Parabéns
Parabéns! Você criou uma arquitetura de memória de agente de IA de longo prazo de duas camadas combinando o Memorystore para Valkey e o AlloyDB AI.
O que você aprendeu
- Implementamos uma arquitetura de memória de duas camadas que separa o estado de sessão ativa de curto prazo dos fatos persistentes de longo prazo.
- Conseguimos uma redução significativa no tamanho do comando ativo e economia no uso total de tokens em relação ao preenchimento de contexto simples, sem perda de acurácia.
- Incorporações automáticas transacionais (
ai.initialize_embeddings) no nível do banco de dados da IA do AlloyDB configuradas. - Realizamos uma pesquisa híbrida de Reciprocal Rank Fusion (
ai.hybrid_search) nativa, combinando a similaridade de vetores (<=>) com a pesquisa de texto completo do PostgreSQL (tsvector). - Criou um extrator de entidades em segundo plano fora da linha de execução (
AsyncMemoryWorker), um avaliador de permissão de execução de ferramentas de três níveis e um mecanismo de compactação de memória de banco de dados. - Anexou o sistema de memória em camadas a um agente autônomo usando o Kit de Desenvolvimento de Agente (ADK) do Google para aplicar restrições de ferramentas e fornecer memória ambiente.
Próximas etapas e referências
- Leia a documentação do AlloyDB AI.
- Saiba como executar a pesquisa vetorial híbrida no AlloyDB.
- Saiba mais sobre o Memorystore para Valkey.