Skip to content

Commit da7ab85

Browse files
authored
Merge pull request #2056 from release/2026.05.1
Release 2026.05.1
2 parents 0346bf1 + 338217b commit da7ab85

12 files changed

Lines changed: 348 additions & 100 deletions

File tree

.github/workflows/_deploy.yaml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -134,7 +134,7 @@ jobs:
134134
- name: configure aws credentials
135135
uses: aws-actions/configure-aws-credentials@v6
136136
with:
137-
role-to-assume: arn:aws:iam::917902836630:role/github-oidc-ssm-version20260317121845461900000008
137+
role-to-assume: arn:aws:iam::917902836630:role/github-oidc-ssm-version
138138
role-session-name: OIDC-GHA-session-version
139139
aws-region: us-east-1
140140
- name: Store version in SSM

.github/workflows/_setup.yaml

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -35,22 +35,22 @@ jobs:
3535
ENV_NAME="${{ inputs.env-name }}"
3636
case "$ENV_NAME" in
3737
"dev")
38-
echo 'role=arn:aws:iam::017925157769:role/github-oidc-deploy20260311132710794900000001' >> "$GITHUB_OUTPUT"
38+
echo 'role=arn:aws:iam::017925157769:role/github-oidc-deploy' >> "$GITHUB_OUTPUT"
3939
;;
4040
"test")
4141
echo 'role=arn:aws:iam::641513112151:role/cmiml-test-oidc-github-role' >> "$GITHUB_OUTPUT"
4242
;;
4343
"uat")
44-
echo 'role=arn:aws:iam::641513112151:role/github-oidc-deploy20260312130551192100000004' >> "$GITHUB_OUTPUT"
44+
echo 'role=arn:aws:iam::641513112151:role/github-oidc-deploy' >> "$GITHUB_OUTPUT"
4545
;;
4646
"stage")
4747
echo 'role=arn:aws:iam::641513112151:role/cmiml-stage-oidc-github-role' >> "$GITHUB_OUTPUT"
4848
;;
4949
"prod")
50-
echo 'role=arn:aws:iam::410431445687:role/github-oidc-deploy20260312143902067900000001' >> "$GITHUB_OUTPUT"
50+
echo 'role=arn:aws:iam::410431445687:role/github-oidc-deploy' >> "$GITHUB_OUTPUT"
5151
;;
5252
"prod-dr-us-west-2")
53-
echo 'role=arn:aws:iam::973422231492:role/github-oidc-deploy20260313122244970500000002' >> "$GITHUB_OUTPUT"
53+
echo 'role=arn:aws:iam::973422231492:role/github-oidc-deploy' >> "$GITHUB_OUTPUT"
5454
;;
5555
*)
5656
echo "Bad environment name"
@@ -100,10 +100,10 @@ jobs:
100100
ENV_NAME="${{ inputs.env-name }}"
101101
case "$ENV_NAME" in
102102
"prod-dr-us-west-2")
103-
echo "role=arn:aws:iam::973422231492:role/github-oidc-ecr20260317131325572800000002" >> "$GITHUB_OUTPUT"
103+
echo "role=arn:aws:iam::973422231492:role/github-oidc-ecr-push" >> "$GITHUB_OUTPUT"
104104
;;
105105
*)
106-
echo "role=arn:aws:iam::917902836630:role/github-oidc-ecr-push20260317121845560500000009" >> "$GITHUB_OUTPUT"
106+
echo "role=arn:aws:iam::917902836630:role/github-oidc-ecr-push" >> "$GITHUB_OUTPUT"
107107
;;
108108
esac
109109

pyproject.toml

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@ dependencies = [
1616
"azure-storage-blob==12.25.*",
1717
"bcrypt~=4.3.0",
1818
"boto3~=1.40.60",
19-
"cryptography>=46.0.6",
19+
"cryptography>=46.0.7",
2020
"ddtrace~=3.17.1",
2121
"fastapi==0.124.*",
2222
"fastapi-mail>=1.5.3",
@@ -28,7 +28,7 @@ dependencies = [
2828
"pyjwt>=2.12.0",
2929
"pymongo==4.13.0",
3030
"pyotp~=2.9.0",
31-
"python-multipart>=0.0.22",
31+
"python-multipart>=0.0.26",
3232
"python-slugify==8.0.4",
3333
"redis==5.2.*",
3434
"sentry-sdk~=2.13",
@@ -85,8 +85,8 @@ dev = [
8585
"polyfactory==2.21.0",
8686
"pre-commit>=4.5.1",
8787
"pyld==2.0.4",
88-
"pytest==8.4.0",
89-
"pytest-asyncio==1.1.0",
88+
"pytest>=9.0.3",
89+
"pytest-asyncio>=1.3.0",
9090
"pytest-cov==6.1.1",
9191
"pytest-env==1.1.5",
9292
"pytest-httpx",

src/broker.py

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,13 @@
1111
from infrastructure.logger import logger
1212

1313
broker: AsyncBroker = (
14-
AioPikaBroker(settings.rabbitmq.url)
14+
AioPikaBroker(
15+
settings.rabbitmq.url,
16+
exchange_name="curious",
17+
queue_name="curious",
18+
declare_exchange_kwargs={"durable": True},
19+
declare_queues_kwargs={"durable": True},
20+
)
1521
.with_result_backend(RedisAsyncResultBackend(settings.redis.url))
1622
.with_formatter(JSONFormatter())
1723
)

src/config/cdn.py

Lines changed: 21 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,16 @@ class CDNSettings(BaseModel):
1919
secret_key: str | None = None
2020
access_key: str | None = None
2121

22+
# KMS
23+
bucket_kms_enabled: bool = False
24+
bucket_kms_key_id: str | None = None
25+
26+
bucket_answer_kms_enabled: bool = False
27+
bucket_answer_kms_key_id: str | None = None
28+
29+
bucket_operation_kms_enabled: bool = False
30+
bucket_operation_kms_key_id: str | None = None
31+
2232
## DR settings
2333
# Override the media bucket name for the DR site
2434
bucket_override: str | None = None
@@ -35,7 +45,7 @@ class CDNSettings(BaseModel):
3545
gcp_endpoint_url: str = "https://storage.googleapis.com"
3646

3747
# Custom Object store endpoint URL
38-
# Usually a custom S3 endpoint or GCP, etc.
48+
# Usually this is a custom S3 endpoint or GCP, etc.
3949
# Locally for minio type stores it is in the form http://localhost:9000
4050
endpoint_url: str | None = None
4151

@@ -51,6 +61,16 @@ def validate_settings(self) -> Self:
5161
"""Validate that domain or endpoint is set. Cannot be both"""
5262
if self.domain and (self.endpoint_url or self.storage_address):
5363
raise ValueError("Either domain or endpoint_url must be set, not both.")
64+
65+
if self.bucket_kms_enabled and not self.bucket_kms_key_id:
66+
raise ValueError("bucket_kms_key_id must be set if bucket_kms_enabled is True")
67+
68+
if self.bucket_answer_kms_enabled and not self.bucket_answer_kms_key_id:
69+
raise ValueError("bucket_answer_kms_key_id must be set if bucket_answer_kms_enabled is True")
70+
71+
if self.bucket_operation_kms_enabled and not self.bucket_operation_kms_key_id:
72+
raise ValueError("bucket_operation_kms_key_id must be set if bucket_operation_kms_enabled is True")
73+
5474
return self
5575

5676
@property

src/infrastructure/storage/storage_client.py

Lines changed: 29 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@
77
from typing import Any, BinaryIO
88

99
import boto3
10-
import httpx
10+
import requests
1111
from botocore.config import Config
1212
from botocore.exceptions import ClientError, EndpointConnectionError
1313
from ddtrace.trace import tracer
@@ -176,7 +176,21 @@ async def list_object(self, key: str):
176176

177177
def generate_presigned_post(self, key) -> dict[str, Any]:
178178
# Not needed ThreadPoolExecutor because there is no any IO operation (no API calls to s3)
179-
return self.client.generate_presigned_post(self._get_bucket_name(), key, ExpiresIn=self.config.ttl_signed_urls)
179+
fields = {}
180+
conditions = []
181+
if self.config.kms_enabled:
182+
fields = {
183+
"x-amz-server-side-encryption": "aws:kms",
184+
"x-amz-server-side-encryption-aws-kms-key-id": self.config.kms_key_id,
185+
}
186+
conditions = [
187+
{"x-amz-server-side-encryption": "aws:kms"},
188+
{"x-amz-server-side-encryption-aws-kms-key-id": self.config.kms_key_id},
189+
]
190+
191+
return self.client.generate_presigned_post(
192+
self._get_bucket_name(), key, ExpiresIn=self.config.ttl_signed_urls, Fields=fields, Conditions=conditions
193+
)
180194

181195
def _copy(self, key, storage_from: "StorageClient", key_from: str | None = None) -> int:
182196
key_from = key_from or key
@@ -211,23 +225,21 @@ async def check(self):
211225
logger.info(f'Check bucket "{storage_bucket}" availability.')
212226
key = "mindlogger.txt"
213227

214-
presigned_data = self.generate_presigned_post(storage_bucket, key)
228+
presigned_data = self.generate_presigned_post(key)
215229

216230
logger.info(f"Presigned POST fields are following: {presigned_data['fields'].keys()}")
217-
file = io.BytesIO(b"")
218-
async with httpx.AsyncClient() as client:
219-
try:
220-
response = await client.post(
221-
presigned_data["url"], data=presigned_data["fields"], files={"file": (key, file)}
222-
)
223-
if response.status_code == http.HTTPStatus.NO_CONTENT:
224-
logger.info(f"Bucket {storage_bucket} is available.")
225-
else:
226-
logger.info(response.content)
227-
raise Exception("File upload error")
228-
except httpx.HTTPError as e:
229-
logger.info("File upload error")
230-
raise e
231+
files = {"file": ("test.txt", b"content")}
232+
233+
try:
234+
response = requests.post(presigned_data["url"], data=presigned_data["fields"], files=files)
235+
if response.status_code == http.HTTPStatus.NO_CONTENT:
236+
logger.info(f"Bucket {storage_bucket} is available.")
237+
else:
238+
logger.info(response.content)
239+
response.raise_for_status()
240+
except requests.exceptions.RequestException as e:
241+
logger.info(f"File upload error: {e}")
242+
raise e
231243

232244
def _check_is_bucket_public(self) -> bool:
233245
# Check the bucket policy

src/infrastructure/storage/storage_config.py

Lines changed: 24 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,16 @@
1+
from typing import Self
2+
3+
from pydantic import model_validator
14
from pydantic_settings import BaseSettings, SettingsConfigDict
25

36
from config import CDNSettings
47

58

69
class StorageConfig(BaseSettings):
7-
"""Configuration class for storage service. Convenience layer between application and object storage."""
10+
"""
11+
Configuration class for storage service. Convenience layer between application and object storage.
12+
This class represents that storage settings to connect with a single object store (i.e., S3 bucket)
13+
"""
814

915
# Custom (S3 or other) endpoint URL
1016
endpoint_url: str | None = None
@@ -13,6 +19,9 @@ class StorageConfig(BaseSettings):
1319
secret_key: str | None = None
1420
access_key: str | None = None
1521

22+
kms_enabled: bool = False
23+
kms_key_id: str | None = None
24+
1625
# Public domain to front storage keys without scheme
1726
# TODO Default to null??
1827
domain: str = ""
@@ -27,6 +36,12 @@ class StorageConfig(BaseSettings):
2736

2837
model_config = SettingsConfigDict(extra="ignore")
2938

39+
@model_validator(mode="after")
40+
def validate_settings(self) -> Self:
41+
if self.kms_enabled and not self.kms_key_id:
42+
raise ValueError("kms_key_id must be set if kms_enabled is True")
43+
return self
44+
3045
@classmethod
3146
def generate_media_settings(cls, cdn_settings: CDNSettings) -> "StorageConfig":
3247
return cls(
@@ -39,6 +54,8 @@ def generate_media_settings(cls, cdn_settings: CDNSettings) -> "StorageConfig":
3954
secret_key=cdn_settings.secret_key,
4055
access_key=cdn_settings.access_key,
4156
ttl_signed_urls=cdn_settings.ttl_signed_urls,
57+
kms_enabled=cdn_settings.bucket_kms_enabled,
58+
kms_key_id=cdn_settings.bucket_kms_key_id,
4259
)
4360

4461
@classmethod
@@ -51,6 +68,8 @@ def generate_answer_settings(cls, cdn_settings: CDNSettings) -> "StorageConfig":
5168
secret_key=cdn_settings.secret_key,
5269
access_key=cdn_settings.access_key,
5370
ttl_signed_urls=cdn_settings.ttl_signed_urls,
71+
kms_enabled=cdn_settings.bucket_answer_kms_enabled,
72+
kms_key_id=cdn_settings.bucket_answer_kms_key_id,
5473
)
5574

5675
@classmethod
@@ -63,6 +82,8 @@ def generate_operations_settings(cls, cdn_settings: CDNSettings) -> "StorageConf
6382
secret_key=cdn_settings.secret_key,
6483
access_key=cdn_settings.access_key,
6584
ttl_signed_urls=cdn_settings.ttl_signed_urls,
85+
kms_enabled=cdn_settings.bucket_operation_kms_enabled,
86+
kms_key_id=cdn_settings.bucket_operation_kms_key_id,
6687
)
6788

6889
@classmethod
@@ -75,4 +96,6 @@ def generate_logs_settings(cls, cdn_settings: CDNSettings) -> "StorageConfig":
7596
secret_key=cdn_settings.secret_key,
7697
access_key=cdn_settings.access_key,
7798
ttl_signed_urls=cdn_settings.ttl_signed_urls,
99+
kms_enabled=cdn_settings.bucket_answer_kms_enabled,
100+
kms_key_id=cdn_settings.bucket_answer_kms_key_id,
78101
)

src/infrastructure/storage/tests/__init__.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,3 +7,5 @@
77
MEDIA_OVERRIDE = "this-is-the-media-override"
88
OPERATIONS_OVERRIDE = "this-is-the-operations-override"
99
ANSWER_OVERRIDE = "this-is-the-answer-override"
10+
11+
KMS_KEY_ID = "arn:aws:kms:test-kms-key"

src/infrastructure/storage/tests/conftest.py

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@
1010
ANSWER_BUCKET_NAME,
1111
ANSWER_OVERRIDE,
1212
DOMAIN,
13+
KMS_KEY_ID,
1314
MEDIA_BUCKET_NAME,
1415
MEDIA_OVERRIDE,
1516
OPERATIONS_BUCKET_NAME,
@@ -137,6 +138,22 @@ async def answer_storage_config(normal_storage_settings: Settings) -> StorageCon
137138
return config
138139

139140

141+
@pytest.fixture
142+
async def answer_storage_kms_config(normal_storage_settings: Settings) -> StorageConfig:
143+
"""Settings for storage"""
144+
config = StorageConfig(
145+
endpoint_url=None,
146+
region="us-east-1",
147+
bucket=normal_storage_settings.cdn.bucket_answer,
148+
access_key=normal_storage_settings.cdn.access_key,
149+
secret_key=normal_storage_settings.cdn.secret_key,
150+
kms_enabled=True,
151+
kms_key_id=KMS_KEY_ID,
152+
)
153+
154+
return config
155+
156+
140157
@pytest.fixture
141158
async def operations_storage_config(normal_storage_settings: Settings) -> StorageConfig:
142159
"""Settings for storage"""

src/infrastructure/storage/tests/test_storage_client.py

Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,15 @@ async def answer_storage_client(self, s3_client) -> StorageClient:
2727

2828
return client
2929

30+
@pytest.fixture
31+
async def answer_storage_client_kms(self, s3_client, answer_storage_kms_config) -> StorageClient:
32+
"""Answer storage client with KMS configured"""
33+
34+
client = StorageClient(answer_storage_kms_config, env="test")
35+
client.client = s3_client
36+
37+
return client
38+
3039
@pytest.fixture
3140
async def media_storage_client_other(self, s3_client) -> StorageClient:
3241
"""Regular storage client with local type settings"""
@@ -66,6 +75,23 @@ async def test_generate_presigned_post(self, answer_storage_client):
6675
assert ANSWER_BUCKET_NAME in data["url"]
6776
assert ANSWER_OVERRIDE not in data["url"]
6877

78+
assert "x-amz-server-side-encryption" not in data["fields"]
79+
assert "x-amz-server-side-encryption-aws-kms-key-id" not in data["fields"]
80+
81+
async def test_generate_presigned_post_kms(self, answer_storage_client_kms):
82+
data = answer_storage_client_kms.generate_presigned_post(FILE_KEY)
83+
assert data is not None
84+
assert ANSWER_BUCKET_NAME in data["url"]
85+
assert ANSWER_OVERRIDE not in data["url"]
86+
87+
assert "x-amz-server-side-encryption" in data["fields"]
88+
assert "x-amz-server-side-encryption-aws-kms-key-id" in data["fields"]
89+
90+
assert data["fields"]["x-amz-server-side-encryption"] == "aws:kms"
91+
assert (
92+
data["fields"]["x-amz-server-side-encryption-aws-kms-key-id"] == answer_storage_client_kms.config.kms_key_id
93+
)
94+
6995
async def test_generate_presigned_post_dr(self, answer_storage_client_dr):
7096
data = answer_storage_client_dr.generate_presigned_post(FILE_KEY)
7197
assert data is not None
@@ -101,3 +127,8 @@ async def test_get_public_url_other_storage_address(self, media_storage_client_o
101127
assert DOMAIN not in url
102128
assert STORAGE_ADDRESS in url
103129
assert FILE_KEY in url
130+
131+
@pytest.mark.usefixtures("answer_bucket")
132+
async def test_check(self, answer_storage_client, s3_client):
133+
answer_storage_client.client = s3_client
134+
await answer_storage_client.check()

0 commit comments

Comments
 (0)