From f340ce643a8b7deee46274cf1592209c120b1b45 Mon Sep 17 00:00:00 2001 From: Leo Iorio Date: Sat, 15 Aug 2026 10:20:35 -0300 Subject: [PATCH] feat(infra): adiciona dynamo_keys e alinha tabela CDK com pk/sk MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Prepara a single-table para USER/MEMBER/SUBSCRIBER e o GSI UserEmailIndex, alinhado ao template, sem reescrever o repositório Dynamo ainda. --- iac/iac/template_dynamo_table.py | 51 +++++++--- iac/iac/template_stack.py | 6 +- src/shared/environments.py | 12 +-- .../infra/external/dynamo/dynamo_keys.py | 96 +++++++++++++++++++ .../repositories/load_user_mock_to_dynamo.py | 54 ++++++++--- .../repositories/user_repository_dynamo.py | 13 ++- .../infra/external/dynamo/test_dynamo_keys.py | 53 ++++++++++ 7 files changed, 244 insertions(+), 41 deletions(-) create mode 100644 src/shared/infra/external/dynamo/dynamo_keys.py create mode 100644 tests/shared/infra/external/dynamo/test_dynamo_keys.py diff --git a/iac/iac/template_dynamo_table.py b/iac/iac/template_dynamo_table.py index 9bc82dc..9b274c9 100644 --- a/iac/iac/template_dynamo_table.py +++ b/iac/iac/template_dynamo_table.py @@ -1,30 +1,55 @@ -from decimal import Decimal from aws_cdk import ( - aws_dynamodb as dynamodb, RemovalPolicy, + RemovalPolicy, + aws_dynamodb as dynamodb, ) from constructs import Construct +# Mantenha alinhado com src.shared.infra.external.dynamo.dynamo_keys / Environments +_EPEP_TABLE_PREFIX = "EpepTable" +_USER_EMAIL_INDEX_NAME = "UserEmailIndex" + +RETAINED_STAGES = {"prod", "homolog"} + class TemplateDynamoTable(Construct): table: dynamodb.Table - def __init__(self, scope: Construct, construct_id: str, **kwargs) -> None: + def __init__(self, scope: Construct, construct_id: str, stage: str = "DEV", **kwargs) -> None: super().__init__(scope, construct_id, **kwargs) + stage_lower = stage.lower() + removal_policy = ( + RemovalPolicy.RETAIN if stage_lower in RETAINED_STAGES else RemovalPolicy.DESTROY + ) + self.table = dynamodb.Table( - self, "TemplateDynamoTable", + self, + "EpepDynamoTable", partition_key=dynamodb.Attribute( - name="PK", - type=dynamodb.AttributeType.STRING + name="pk", + type=dynamodb.AttributeType.STRING, ), sort_key=dynamodb.Attribute( - name="SK", - type=dynamodb.AttributeType.STRING - ), + name="sk", + type=dynamodb.AttributeType.STRING, + ), billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST, - removal_policy=RemovalPolicy.DESTROY + removal_policy=removal_policy, + table_name=f"{_EPEP_TABLE_PREFIX}-{stage_lower}", + point_in_time_recovery_specification=dynamodb.PointInTimeRecoverySpecification( + point_in_time_recovery_enabled=(stage_lower == "prod") + ), ) - - - + self.table.add_global_secondary_index( + index_name=_USER_EMAIL_INDEX_NAME, + partition_key=dynamodb.Attribute( + name="gsi2pk", + type=dynamodb.AttributeType.STRING, + ), + sort_key=dynamodb.Attribute( + name="gsi2sk", + type=dynamodb.AttributeType.STRING, + ), + projection_type=dynamodb.ProjectionType.ALL, + ) diff --git a/iac/iac/template_stack.py b/iac/iac/template_stack.py index dbaf9b4..0556518 100644 --- a/iac/iac/template_stack.py +++ b/iac/iac/template_stack.py @@ -34,13 +34,13 @@ def __init__(self, scope: Construct, construct_id: str, **kwargs) -> None: } ) - self.dynamo_table = TemplateDynamoTable(self, "TemplateDynamoTable") + self.dynamo_table = TemplateDynamoTable(self, "TemplateDynamoTable", stage="DEV") ENVIRONMENT_VARIABLES = { "STAGE": "DEV", "DYNAMO_TABLE_NAME": self.dynamo_table.table.table_name, - "DYNAMO_PARTITION_KEY": "PK", - "DYNAMO_SORT_KEY": "SK", + "DYNAMO_PARTITION_KEY": "pk", + "DYNAMO_SORT_KEY": "sk", "REGION": self.region, } diff --git a/src/shared/environments.py b/src/shared/environments.py index ce8c825..0dcb61a 100644 --- a/src/shared/environments.py +++ b/src/shared/environments.py @@ -24,7 +24,7 @@ class Environments: stage: STAGE s3_bucket_name: str region: str - endpoint_url: str = None + dynamo_endpoint_url: str = None # DynamoDB Local (ex: http://localhost:8000); None na AWS dynamo_table_name: str dynamo_partition_key: str dynamo_sort_key: str @@ -46,16 +46,16 @@ def load_envs(self): if self.stage == STAGE.TEST: self.s3_bucket_name = "bucket-test" self.region = "sa-east-1" - self.endpoint_url = "http://localhost:8000" + self.dynamo_endpoint_url = "http://localhost:8000" self.dynamo_table_name = "user_mss_template-table" - self.dynamo_partition_key = "PK" - self.dynamo_sort_key = "SK" + self.dynamo_partition_key = "pk" + self.dynamo_sort_key = "sk" self.cloud_front_distribution_domain = "https://d3q9q9q9q9q9q9.cloudfront.net" else: self.s3_bucket_name = os.environ.get("S3_BUCKET_NAME") self.region = os.environ.get("REGION") - self.endpoint_url = os.environ.get("ENDPOINT_URL") + self.dynamo_endpoint_url = os.environ.get("DYNAMO_ENDPOINT_URL") self.dynamo_table_name = os.environ.get("DYNAMO_TABLE_NAME") self.dynamo_partition_key = os.environ.get("DYNAMO_PARTITION_KEY") self.dynamo_sort_key = os.environ.get("DYNAMO_SORT_KEY") @@ -86,7 +86,7 @@ def get_observability() -> IObservability: def get_envs() -> "Environments": """ Returns the Environments object. This method should be used to get the Environments object instead of instantiating it directly. - :return: Environments (stage={self.stage}, s3_bucket_name={self.s3_bucket_name}, region={self.region}, endpoint_url={self.endpoint_url}) + :return: Environments (stage={self.stage}, region={self.region}, dynamo_table_name={self.dynamo_table_name}, dynamo_endpoint_url={self.dynamo_endpoint_url}) """ envs = Environments() diff --git a/src/shared/infra/external/dynamo/dynamo_keys.py b/src/shared/infra/external/dynamo/dynamo_keys.py new file mode 100644 index 0000000..c4f3d25 --- /dev/null +++ b/src/shared/infra/external/dynamo/dynamo_keys.py @@ -0,0 +1,96 @@ +"""Convenções de chaves Dynamo (tabela base + GSI UserEmailIndex). + +Tabela base (single-table): + pk = USER | MEMBER | SUBSCRIBER + sk = USER# | MEMBER# | SUBSCRIBER# + +GSI2 (UserEmailIndex) — access pattern "user por email": + gsi2pk = EMAIL# + gsi2sk = USER# + +Uso no repository Dynamo (exemplo):: + + from boto3.dynamodb.conditions import Key + from src.shared.infra.external.dynamo.dynamo_keys import ( + GSI2_NAME, GSI2_PK_ATTR, gsi2_partition_key, + ) + + resp = self.dynamo.query( + KeyConditionExpression=Key(GSI2_PK_ATTR).eq(gsi2_partition_key(user_email)), + IndexName=GSI2_NAME, + ) +""" + +from enum import Enum +from typing import Any +from uuid import UUID + +from pydantic import EmailStr + +# Nomes dos atributos da tabela base (alinhados ao CDK: pk/sk) +PK_ATTR = "pk" +SK_ATTR = "sk" + +# GSI2 — access pattern: buscar user por email (denso: email sempre presente) +# Alinhado a iac/iac/template_dynamo_table.py (UserEmailIndex) +GSI2_NAME = "UserEmailIndex" +GSI2_PK_ATTR = "gsi2pk" +GSI2_SK_ATTR = "gsi2sk" + +STORAGE_KEY_ATTRS = ( + PK_ATTR, SK_ATTR, + GSI2_PK_ATTR, GSI2_SK_ATTR, +) + + +class EntityKind(str, Enum): + USER = "USER" + MEMBER = "MEMBER" + SUBSCRIBER = "SUBSCRIBER" + + +def partition_key(kind: EntityKind) -> str: + """PK da tabela base — coleção (se repete para todos os items do kind).""" + return kind.value + + +def sort_key(id: UUID, kind: EntityKind) -> str: + """SK da tabela base — identidade única dentro da coleção.""" + return f"{kind.value}#{id}" + + +def gsi2_partition_key(user_email: EmailStr) -> str: + """ + PK do GSI2 — agrupa por email. + + Ex.: EMAIL#user@example.com + """ + return f"EMAIL#{user_email}" + + +def gsi2_sort_key(user_id: UUID) -> str: + """ + SK do GSI2 — identidade do user no índice de email. + + Ex.: USER# + """ + return sort_key(id=user_id, kind=EntityKind.USER) + + +def build_gsi2_attributes( + user_email: EmailStr, + user_id: UUID, +) -> dict[str, str]: + """ + GSI denso: email é obrigatório na entidade User, então todo user + recebe gsi2pk/gsi2sk e entra no UserEmailIndex. + """ + return { + GSI2_PK_ATTR: gsi2_partition_key(user_email=user_email), + GSI2_SK_ATTR: gsi2_sort_key(user_id=user_id), + } + + +def strip_keys(item: dict[str, Any]) -> dict[str, Any]: + """Remove atributos de storage (pk/sk/gsi) antes do model_validate.""" + return {k: v for k, v in item.items() if k not in STORAGE_KEY_ATTRS} diff --git a/src/shared/infra/repositories/load_user_mock_to_dynamo.py b/src/shared/infra/repositories/load_user_mock_to_dynamo.py index 8c55305..2ad0852 100644 --- a/src/shared/infra/repositories/load_user_mock_to_dynamo.py +++ b/src/shared/infra/repositories/load_user_mock_to_dynamo.py @@ -5,11 +5,19 @@ from src.shared.infra.repositories.user_repository_dynamo import UserRepositoryDynamo from src.shared.infra.repositories.user_repository_mock import UserRepositoryMock from src.shared.environments import Environments +from src.shared.infra.external.dynamo.dynamo_keys import ( + GSI2_NAME, + GSI2_PK_ATTR, + GSI2_SK_ATTR, +) def setup_dynamo_table(): - dynamo_table_name = "user_mss_template-table" - endpoint_url = "http://localhost:8000" + envs = Environments.get_envs() + dynamo_table_name = envs.dynamo_table_name + endpoint_url = envs.dynamo_endpoint_url + pk = envs.dynamo_partition_key + sk = envs.dynamo_sort_key print("Setting up DynamoDB table...") dynamo_client = boto3.client('dynamodb', endpoint_url=endpoint_url) @@ -22,24 +30,41 @@ def setup_dynamo_table(): TableName=dynamo_table_name, KeySchema=[ { - 'AttributeName': 'PK', + 'AttributeName': pk, 'KeyType': 'HASH' }, { - 'AttributeName': 'SK', + 'AttributeName': sk, 'KeyType': 'RANGE' } ], AttributeDefinitions=[ { - 'AttributeName': 'PK', + 'AttributeName': pk, 'AttributeType': 'S' }, { - 'AttributeName': 'SK', + 'AttributeName': sk, 'AttributeType': 'S' - } - + }, + { + 'AttributeName': GSI2_PK_ATTR, + 'AttributeType': 'S' + }, + { + 'AttributeName': GSI2_SK_ATTR, + 'AttributeType': 'S' + }, + ], + GlobalSecondaryIndexes=[ + { + "IndexName": GSI2_NAME, + "KeySchema": [ + {"AttributeName": GSI2_PK_ATTR, "KeyType": "HASH"}, + {"AttributeName": GSI2_SK_ATTR, "KeyType": "RANGE"}, + ], + "Projection": {"ProjectionType": "ALL"}, + }, ], BillingMode='PAY_PER_REQUEST', ) @@ -56,13 +81,13 @@ def setup_dynamo_table(): table.put_item( Item={ - 'PK': 'COUNTER', - 'SK': 'COUNTER', + pk: 'COUNTER', + sk: 'COUNTER', 'COUNTER': Decimal(0) } ) - print('Table "user_mss_template-table" created!') + print(f'Table "{dynamo_table_name}" created!') else: print("Table already exists!") @@ -89,15 +114,16 @@ def load_mock_to_real_dynamo(): count = 0 + envs = Environments.get_envs() dynamodb = boto3.resource('dynamodb') - table = dynamodb.Table(dynamo_table_name=Environments.get_envs().dynamo_table_name) + table = dynamodb.Table(envs.dynamo_table_name) print("Adding counter to table") table.put_item( Item={ - 'PK': 'COUNTER', - 'SK': 'COUNTER', + envs.dynamo_partition_key: 'COUNTER', + envs.dynamo_sort_key: 'COUNTER', 'COUNTER': Decimal(0) } ) diff --git a/src/shared/infra/repositories/user_repository_dynamo.py b/src/shared/infra/repositories/user_repository_dynamo.py index 15accd9..17408c7 100644 --- a/src/shared/infra/repositories/user_repository_dynamo.py +++ b/src/shared/infra/repositories/user_repository_dynamo.py @@ -20,11 +20,14 @@ def sort_key_format(user_id: int) -> str: return f"#{user_id}" def __init__(self): - self.dynamo = DynamoDatasource(endpoint_url=Environments.get_envs().endpoint_url, - dynamo_table_name=Environments.get_envs().dynamo_table_name, - region=Environments.get_envs().region, - partition_key=Environments.get_envs().dynamo_partition_key, - sort_key=Environments.get_envs().dynamo_sort_key) + envs = Environments.get_envs() + self.dynamo = DynamoDatasource( + endpoint_url=envs.dynamo_endpoint_url, + dynamo_table_name=envs.dynamo_table_name, + region=envs.region, + partition_key=envs.dynamo_partition_key, + sort_key=envs.dynamo_sort_key, + ) def get_user(self, user_id: int) -> User: resp = self.dynamo.get_item(partition_key=self.partition_key_format(user_id), sort_key=self.sort_key_format(user_id)) diff --git a/tests/shared/infra/external/dynamo/test_dynamo_keys.py b/tests/shared/infra/external/dynamo/test_dynamo_keys.py new file mode 100644 index 0000000..88752cc --- /dev/null +++ b/tests/shared/infra/external/dynamo/test_dynamo_keys.py @@ -0,0 +1,53 @@ +from uuid import UUID + +from src.shared.infra.external.dynamo.dynamo_keys import ( + EntityKind, + GSI2_NAME, + GSI2_PK_ATTR, + GSI2_SK_ATTR, + PK_ATTR, + SK_ATTR, + build_gsi2_attributes, + partition_key, + sort_key, + strip_keys, +) + + +class Test_DynamoKeys: + def test_user_base_keys(self): + user_id = UUID("11111111-1111-1111-1111-111111111111") + + assert PK_ATTR == "pk" + assert SK_ATTR == "sk" + assert partition_key(EntityKind.USER) == "USER" + assert sort_key(user_id, EntityKind.USER) == f"USER#{user_id}" + + def test_member_and_subscriber_kinds(self): + member_id = UUID("22222222-2222-2222-2222-222222222222") + subscriber_id = UUID("33333333-3333-3333-3333-333333333333") + + assert partition_key(EntityKind.MEMBER) == "MEMBER" + assert sort_key(member_id, EntityKind.MEMBER) == f"MEMBER#{member_id}" + assert partition_key(EntityKind.SUBSCRIBER) == "SUBSCRIBER" + assert sort_key(subscriber_id, EntityKind.SUBSCRIBER) == f"SUBSCRIBER#{subscriber_id}" + + def test_user_email_gsi(self): + user_id = UUID("11111111-1111-1111-1111-111111111111") + attrs = build_gsi2_attributes("user@example.com", user_id) + + assert GSI2_NAME == "UserEmailIndex" + assert attrs[GSI2_PK_ATTR] == "EMAIL#user@example.com" + assert attrs[GSI2_SK_ATTR] == f"USER#{user_id}" + + def test_strip_keys(self): + user_id = UUID("11111111-1111-1111-1111-111111111111") + item = { + "pk": "USER", + "sk": f"USER#{user_id}", + "gsi2pk": "EMAIL#user@example.com", + "gsi2sk": f"USER#{user_id}", + "email": "user@example.com", + } + + assert strip_keys(item) == {"email": "user@example.com"}