Pular para conteúdo

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.