DynamoDB: entidades, chaves e repositories¶
Modelo single-table¶
AgentTable, em template.yaml, define uma tabela gogenetic-agent-${Environment}, capacidade PAY_PER_REQUEST, chave composta PK/SK e índice GSI1 com projeção ALL. O atributo de expiração é ttl. O nome efetivo consumido pelo código vem de DYNAMODB_TABLE; o factory recebe a tabela ou a cria como referência via boto3.resource("dynamodb").Table(...).
Single-table significa guardar diferentes tipos de registros na mesma tabela com prefixos de chave. Não significa que todos têm o mesmo formato, TTL ou atomicidade. src/core/integrations/dynamodb/keys.py centraliza as chaves principais; crm_sync.py define sua própria chave de projeção.
Inventário de itens¶
| Entidade | PK | SK | Índice / retenção |
|---|---|---|---|
| Sessão | SESSION#{phone} |
SESSION#{session_id} |
GSI1PK=STATE#{state}, GSI1SK=created_at ISO; sem TTL automático |
| Pointer ativo | SESSION#{phone} |
ACTIVE |
Guarda session_id, updated_at; sem GSI |
| Mensagem | SESSION#{phone} |
MSG#{timestamp_iso}#{message_id} |
Histórico por telefone/sessão; sem TTL |
| Demanda fora do portfólio | DEMAND#OOP |
DEMAND#{timestamp_iso}#{demand_id} |
Consulta cronológica; sem TTL |
| Receipt inbound | DEBOUNCE#{phone} |
INBOUND#{sha256(message_id)} |
TTL recebido+7 dias; flags de roteamento/despacho/conclusão |
| Última entrada | DEBOUNCE#{phone} |
LATEST_INBOUND |
latest, lista pending, marcadores de commit/invalidação; sem TTL |
| Buffer legado | DEBOUNCE#{phone} |
BUFFER |
ttl=first_message_at+30s |
| Resposta pendente legada | DEBOUNCE#{phone} |
PENDING_REPLY |
Texto, sessão, estado e primeira mensagem; sem TTL no serializer |
| Projeção comercial | CRM_SYNC#{session_id} |
META |
GSI1PK=CRM_SYNC#{status}, GSI1SK=next_attempt_at#{session_id} |
| Recurso administrativo atual | ADMIN#RESOURCE#{resource_type} |
RESOURCE#{resource_id} |
Prompt/package atual, conteúdo, versão, ativo e autoria |
| Versão de recurso | ADMIN#RESOURCE#{resource_type}#{resource_id} |
VERSION#{version:06d} |
Histórico de versões |
| Configuração operacional | ADMIN#CONFIG |
CONFIG#{key} |
Valor, autor e timestamp |
| Auditoria administrativa | ADMIN#AUDIT |
AUDIT#{timestamp_iso}#{event_id} |
Eventos de mudança |
| Consentimento | ADMIN#CONSENT |
CONSENT#{timestamp_iso}#{phone}#{consent_id} |
Registro próprio, não uma aprovação clínica/jurídica automática |
| Grant de acesso operador | OPERATOR#AUTH |
LOGIN#{token_hash} |
TTL, telefone; token bruto não é a chave |
| Sessão de operador | OPERATOR#AUTH |
SESSION#{token_hash} |
TTL, telefone, CSRF |
A disponibilidade humana é um OperationalConfig com chave operator_availability: mapa de telefone para timestamp de disponivel. Não é uma linha por operador com TTL DynamoDB. A aplicação exige mesmo dia local e idade dentro de HANDOFF_OPERATOR_AVAILABILITY_TTL_HOURS (default 24). Grants de login duram no máximo 10 minutos ou até a próxima meia-noite local; a sessão do operador expira à meia-noite. O código verifica validade, sem depender apenas da remoção física por TTL.
Sessões¶
src/core/integrations/dynamodb/sessions.py define SessionRepository e DynamoSessionRepository.
| Operação | Entrada → retorno | Uso e condições |
|---|---|---|
get_active(phone) |
telefone → Session | None |
Lê pointer e sessão com ConsistentRead=True; adiciona IDs/status da projeção CRM |
get(phone, session_id) |
identidade → Session | None |
Recupera sessão específica com leitura consistente |
put(session) |
entidade → None |
Transação de sessão+ACTIVE; pode incluir intenção CRM; não é CAS de ownership |
try_put_agent_result(session, expected_latest_message_id, expected_active_session_id=None, replaced_sessions=(), require_latest=True) |
snapshot → bool |
Persiste resultado só se latest/pointer/estado ainda permitem; usada pelo TurnGuard |
transition_state(session, expected_state, new_state, extra_updates=None, expected_assigned_to=None, require_no_live_human_send_at=None) |
transição → None |
Condiciona estado, opcionalmente operador e ausência de lock vivo; com CRM também compara updated_at |
list_by_state(state) |
enum → list[Session] |
Query paginada GSI1, ordenada por criação; usada pelo inbox e jobs |
acquire_human_send_lock(session, lock_id, now, expires_at) |
lock → None |
Protege envio manual concorrente e ownership |
release_human_send_lock(session, lock_id) |
lock → bool |
Remove somente o lock com ID correspondente |
mark_handoff_renotified(session, notified_at) |
sessão/hora → None |
Atualização condicional, contador e timestamp de renotificação |
mark_handoff_customer_renotified(session, notified_at) |
sessão/hora → None |
Guarda confirmação de renotificação ao cliente |
set_active_pointer(phone, session_id) |
identidade → None |
Troca explícita do pointer; exige cuidado fora de transações completas |
_transact() prepara a intenção comercial na mesma transação quando crm_sync_enabled=True. Repete até oito tentativas somente para conflito CAS da projeção, sem ignorar falha das condições da sessão. Conflitos de estado mapeiam para DynamoConditionFailedError; indisponibilidade/autorização do SDK não vira sucesso.
Validação e atomicidade são camadas diferentes
Session.transition_to() valida a máquina de estados em memória. transition_state() do repository verifica condições de banco, mas não consulta LEGAL_TRANSITIONS. Um consumidor novo precisa aplicar a regra de domínio e a condição atômica.
Mensagens e demandas¶
src/core/integrations/dynamodb/messages.py contém MessageRepository/DynamoMessageRepository. append(Message) grava um item; a chave com timestamp e ID permite repetir o mesmo registro sem criar outra chave, mas não substitui a deduplicação de provider. append_inbound_if_session_state(message, expected_state) retorna bool e transaciona criação exclusiva da mensagem, estado esperado e pointer ativo: é a proteção do inbox humano.
history(phone, limit) consulta mensagens do telefone. history_for_session(phone, session_id, limit) pagina e filtra para obter o histórico correto da conversa, retornando ordem cronológica. O invocador usa a segunda. list_recent_llm(limit) atende custos/observabilidade e faz scan; não é consulta indexada dedicada. mark_superseded(message) marca erro superseded para retirar uma resposta antiga de replay/histórico do agente.
src/core/integrations/dynamodb/demands.py contém DemandRepository/DynamoDemandRepository: record(OutOfPortfolioDemand) → None e list_recent(limit) → list[OutOfPortfolioDemand]. O item inclui serviço solicitado, mensagem original, segmento estimado e instância. O ID de demanda é novo; o registro não possui condição exclusiva baseada em provider ID. Em corridas/retries, não presuma deduplicação de negócio apenas por usar DynamoDB.
Receipts e debounce¶
src/core/integrations/dynamodb/debounce.py contém DebounceRepository, DynamoDebounceRepository e InboundRegistration(is_duplicate, already_dispatched, already_completed).
| Operações | Retorno / garantia |
|---|---|
register_inbound(InboundTurn) |
InboundRegistration; receipt exclusivo e atualização atômica de latest/pending |
mark_agent_routed(turn) |
bool; não expõe ao bot receipt claimed síncrono ou concluído |
mark_dispatched(turn) |
None; requer receipt existente |
try_claim_sync(turn) / release_sync_claim(turn) |
bool / None; claim exclusivo de rota síncrona, liberável numa corrida de roteamento |
is_latest(phone, message_id) / is_invalidated(phone) |
bool; leitura consistente do marcador |
pending_inbounds(phone) |
list[InboundTurn] em ordem de registro, só bot-owned |
invalidate_pending(...) |
Barreira que impede turnos antigos após encerramento de debug |
try_commit(phone, message_id) |
bool; valida latest e sessão/estado antes da ação |
is_completed(turn) / mark_completed(turn) |
bool / None; conclusão durável e remoção condicional da lista pending |
get(phone) / put(buffer) / delete(phone) |
Buffer legado; grava e lê DebounceBuffer |
delete_if_first_message_at(phone, first_message_at) |
Só apaga o buffer esperado |
get_pending_outbound(phone) / put_pending_outbound(pending) / delete_pending_outbound(phone) |
Compatibilidade de replay legado |
try_claim_sync assume execução no máximo uma vez após o claim: uma queda entre efeito e marca de conclusão não permite distinguir automaticamente “não enviado” de “enviado sem persistir confirmação”. A rota de ACK humano trata falha de envio explicitamente, liberando o claim para retry.
Projeção comercial¶
src/core/integrations/dynamodb/crm_sync.py contém CrmSyncRepository. get(session_id) retorna um mapa ou None; prepare(session) retorna (row, etag_anterior) ou None sem escrever diretamente; a sessão incorpora esse intent em transação. change(session_id, mutate) usa CAS por etag, até oito tentativas, retornando o mapa atualizado ou None. Esgotamento lança RuntimeError("crm_sync_concurrent_update"). list_status(status) pagina GSI1; problems() agrega PENDING, BLOCKED, UNCERTAIN.
desired_projection() cria contato/deal e stage lógico. A revisão é SHA-256 do JSON ordenado, permitindo não reenviar conteúdo idêntico. IDs externos e marcadores como creation_started sobrevivem a substituições do snapshot de sessão. NEW sem deal existente ainda não publica uma projeção: a alocação ocorre antes do guard definitivo do primeiro turno.
Administração e autenticação¶
Todas as classes abaixo ficam em src/core/integrations/dynamodb/admin.py, exceto autenticação.
| Classe | Operações / retorno | Condições e erros |
|---|---|---|
DynamoVersionedResourceRepository |
list_current(type), get_current(type,id), list_versions(type,id), put_next_version(...), deactivate(...), rollback(...,target_version) → recursos |
Versão e current são dois put_item separados; não há CAS para edições simultâneas. Rollback ausente lança ValueError e rollback cria uma versão nova |
DynamoOperationalConfigRepository |
list_all(), get(key), set(OperationalConfig) |
Set genérico sem CAS; disponibilidade usa helper próprio de atualização atômica por updated_at |
DynamoAdminAuditRepository |
record(AdminAuditEvent), list_recent(limit) |
Append e consulta cronológica; registra autor/ação/recurso |
DynamoConsentRepository |
record(ConsentRecord), list_recent(limit) |
Armazena evento de consentimento; sem expiração genérica |
DynamoOperatorAuthRepository em operator_auth.py |
put_login_grant(token_hash,operator_phone,expires_at), consume_login_grant(token_hash,now), put_session(token_hash,session), get_session(token_hash), delete_session(token_hash) |
Criação exclusiva; consumo de grant por delete condicional existente e ttl > now, retorna grant ou None; sessão vencida não autentica |
Não aplique as garantias CAS das sessões a todos os repositories. Recursos administrativos e consultas list_current/list_versions/list_all/list_recent possuem implementações próprias; várias fazem uma única query, portanto não se deve assumir paginação ilimitada em listas administrativas extensas.
Exemplos sintéticos de itens¶
{
"PK":"SESSION#<TELEFONE_FICTICIO>",
"SK":"SESSION#sessao-exemplo",
"session_id":"sessao-exemplo",
"phone_number":"<TELEFONE_FICTICIO>",
"whatsapp_instance":"gogenetic",
"state":"QUALIFYING",
"current_agent":"qualificacao",
"segment":"PESQUISA",
"intent":"servico",
"qualification_data":{"tipo_amostra":"amostra sintética"},
"cadastral_data":{},
"created_at":"2026-09-14T12:00:00+00:00",
"updated_at":"2026-09-14T12:00:05+00:00",
"GSI1PK":"STATE#QUALIFYING",
"GSI1SK":"2026-09-14T12:00:00+00:00"
}
O exemplo omite campos opcionais/defaults para leitura. O contrato completo de sessão e serialization.py definem o formato efetivo. Datas viram ISO; enums viram strings; dados numéricos DynamoDB usam Decimal e são normalizados pelos serializers. O GSI também contém projeções comerciais: filtrar somente por existência de GSI1PK não identifica uma sessão.
Testar, diagnosticar e alterar¶
Use tests/unit/integrations/test_dynamodb_repos.py, test_operator_auth_repository.py, tests/unit/application/test_crm_sync.py, test_operator_availability.py e test_human_conversation.py. Testes com Moto verificam regras sem dados remotos. Para campos novos, revise entidade, ida/volta de serialization.py, compatibilidade de linhas antigas, índices e consumers da UI. Para nova chave, centralize o builder e teste leitura após escrita e concorrência pertinente.
TTL habilitado não constitui política de retenção completa: sessões, mensagens, demandas e recursos administrativos permanecem sem expiração automática neste código. Qualquer política de descarte/privacidade precisa tratar esses registros explicitamente. Veja produção e observabilidade.