Compare commits

...

2 Commits

6 changed files with 79 additions and 51 deletions

View File

@@ -1,7 +0,0 @@
BEGIN;
ALTER TABLE sbp_withdrawals
DROP COLUMN IF EXISTS wallet_address,
DROP COLUMN IF EXISTS sender_wallet_address;
COMMIT;

View File

@@ -12,10 +12,10 @@ services:
APP_HOST: "0.0.0.0"
APP_PORT: "8000"
APP_WORKERS: "2"
KEYDB_REMOTE_HOST: "pay_keydb"
KEYDB_REMOTE_PORT: "${KEYDB_PORT:-6380}"
KEYDB_REMOTE_PASSWORD: "${REDIS_PASSWORD}"
KEYDB_REMOTE_DB: "${REDIS_DB:-0}"
KEYDB_HOST: "pay_keydb"
KEYDB_PORT: "${KEYDB_PORT:-6380}"
KEYDB_PASSWORD: "${REDIS_PASSWORD}"
KEYDB_DB: "${REDIS_DB:-0}"
env_file:
- .env
depends_on:

View File

@@ -4,7 +4,7 @@ from src.infrastructure.config import settings
def create_redis_client(url:str|None=None) -> Redis:
redis_url = url or settings.KEYDB_REMOTE_URL
redis_url = url or settings.KEYDB_MARKET_URL
return redis.from_url(
redis_url,
max_connections=50,

View File

@@ -1,14 +1,14 @@
from __future__ import annotations
import os
from functools import lru_cache
from typing import List,Literal
import os
from dotenv import load_dotenv, find_dotenv
from dotenv import find_dotenv,load_dotenv
from pydantic import AliasChoices,Field,field_validator,model_validator
from pydantic_settings import BaseSettings,SettingsConfigDict
from src.infrastructure.vault import create_hvac_client_from_approle, read_kv2_secret
env_file = find_dotenv(".env")
env_file = find_dotenv('.env')
if env_file:
load_dotenv(env_file)
@@ -80,6 +80,10 @@ class Settings(BaseSettings):
KEYDB_REMOTE_PORT: int | None = None
KEYDB_REMOTE_PASSWORD: str | None = None
KEYDB_REMOTE_DB: int | None = None
KEYDB_HOST: str | None = Field(default=None,validation_alias=AliasChoices('KEYDB_HOST','KEYDB_CACHE_HOST'))
KEYDB_PORT: int | None = Field(default=None,validation_alias=AliasChoices('KEYDB_PORT','KEYDB_CACHE_PORT'))
KEYDB_PASSWORD: str | None = Field(default=None,validation_alias=AliasChoices('KEYDB_PASSWORD','KEYDB_CACHE_PASSWORD'))
KEYDB_DB: int | None = Field(default=None,validation_alias=AliasChoices('KEYDB_DB','KEYDB_CACHE_DB'))
RABBIT_HOST: str = "localhost"
RABBIT_PORT: int = 5672
@@ -128,7 +132,7 @@ class Settings(BaseSettings):
return v
return normalize_vault_base_url(v)
@model_validator(mode="before")
@model_validator(mode='before')
@classmethod
def load_from_vault(cls, data: dict):
if not isinstance(data, dict):
@@ -284,19 +288,6 @@ class Settings(BaseSettings):
else:
data['KEYDB_REMOTE_PASSWORD'] = None
keydb_host_override = os.getenv('KEYDB_REMOTE_HOST') or data.get('KEYDB_REMOTE_HOST')
keydb_port_override = os.getenv('KEYDB_REMOTE_PORT') or data.get('KEYDB_REMOTE_PORT')
keydb_db_override = os.getenv('KEYDB_REMOTE_DB') or data.get('KEYDB_REMOTE_DB')
keydb_password_override = os.getenv('KEYDB_REMOTE_PASSWORD') or data.get('KEYDB_REMOTE_PASSWORD')
if keydb_host_override is not None and str(keydb_host_override).strip():
data['KEYDB_REMOTE_HOST'] = str(keydb_host_override).strip()
if keydb_port_override is not None and str(keydb_port_override).strip():
data['KEYDB_REMOTE_PORT'] = int(keydb_port_override)
if keydb_db_override is not None and str(keydb_db_override).strip():
data['KEYDB_REMOTE_DB'] = int(keydb_db_override)
if keydb_password_override is not None and str(keydb_password_override).strip():
data['KEYDB_REMOTE_PASSWORD'] = str(keydb_password_override).strip()
itpay_public_id = data.get('ITPAY_PUBLIC_ID') or os.getenv('ITPAY_PUBLIC_ID')
itpay_api_secret = data.get('ITPAY_API_SECRET') or os.getenv('ITPAY_API_SECRET')
if itpay_public_id is not None and str(itpay_public_id).strip() and itpay_api_secret is not None and str(itpay_api_secret).strip():
@@ -348,20 +339,41 @@ class Settings(BaseSettings):
@property
def REDIS_URL(self) -> str:
return self.KEYDB_REMOTE_URL
return self.KEYDB_MARKET_URL
@staticmethod
def _redis_url(*, host: str, port: int, password: str | None, db: int) -> str:
auth = f':{password}@' if password else ''
return f'redis://{auth}{host}:{port}/{db}'
@property
def KEYDB_MARKET_URL(self) -> str:
if self.KEYDB_REMOTE_HOST is None or self.KEYDB_REMOTE_PORT is None or self.KEYDB_REMOTE_DB is None:
raise RuntimeError('Vault KeyDB settings are required for market data')
return self._redis_url(
host=self.KEYDB_REMOTE_HOST,
port=int(self.KEYDB_REMOTE_PORT),
password=self.KEYDB_REMOTE_PASSWORD,
db=int(self.KEYDB_REMOTE_DB),
)
@property
def KEYDB_REMOTE_URL(self) -> str:
host = self.KEYDB_REMOTE_HOST or self.REDIS_HOST
port = int(self.KEYDB_REMOTE_PORT) if self.KEYDB_REMOTE_PORT is not None else int(self.REDIS_PORT)
password = self.KEYDB_REMOTE_PASSWORD if self.KEYDB_REMOTE_PASSWORD is not None else self.REDIS_PASSWORD
db = int(self.KEYDB_REMOTE_DB) if self.KEYDB_REMOTE_DB is not None else int(self.REDIS_DB)
return self._redis_url(host=host, port=port, password=password, db=db)
return self.KEYDB_MARKET_URL
@property
def KEYDB_CACHE_URL(self) -> str | None:
if self.KEYDB_HOST is None or not self.KEYDB_HOST.strip():
return None
if self.KEYDB_PORT is None:
return None
db = int(self.KEYDB_DB) if self.KEYDB_DB is not None else 0
return self._redis_url(
host=self.KEYDB_HOST.strip(),
port=int(self.KEYDB_PORT),
password=self.KEYDB_PASSWORD,
db=db,
)
@property
def RABBIT_URL(self) -> str:

View File

@@ -40,16 +40,33 @@ async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
logger.set_instance_id(instance_id)
logger.info(f'Users service instance started with id {instance_id}')
app.state.redis_remote = create_redis_client(settings.KEYDB_REMOTE_URL)
app.state.redis = app.state.redis_remote
app.state.redis_market = create_redis_client(settings.KEYDB_MARKET_URL)
app.state.redis_remote = app.state.redis_market
app.state.redis_cache = None
app.state.redis = None
try:
await app.state.redis_remote.ping()
logger.info('KeyDB connection established')
await app.state.redis_market.ping()
logger.info('Market KeyDB connection established')
except Exception as exception:
logger.error(f'KeyDB connection failed: {exception}')
await app.state.redis_remote.aclose()
logger.error(f'Market KeyDB connection failed: {exception}')
await app.state.redis_market.aclose()
raise
cache_url = settings.KEYDB_CACHE_URL
if cache_url:
app.state.redis_cache = create_redis_client(cache_url)
app.state.redis = app.state.redis_cache
try:
await app.state.redis_cache.ping()
logger.info('Cache KeyDB connection established')
except Exception as exception:
logger.error(f'Cache KeyDB connection failed: {exception}')
await app.state.redis_cache.aclose()
await app.state.redis_market.aclose()
raise
else:
logger.info('Cache KeyDB is not configured')
jwt_store = JwtKeyStore(
vault_addr=settings.VAULT_ADDR,
vault_role_id=settings.VAULT_ROLE_ID,
@@ -72,7 +89,10 @@ async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
sched = getattr(app.state,'jwt_keys_scheduler',None)
if sched:
sched.shutdown(wait=False)
await app.state.redis_remote.aclose()
redis_cache = getattr(app.state,'redis_cache',None)
if redis_cache is not None:
await redis_cache.aclose()
await app.state.redis_market.aclose()
logger.info(f'Pay service instance ended with id {instance_id}')

View File

@@ -9,19 +9,22 @@ from src.infrastructure.cache import KeydbCache,RemoteCache
def get_redis_remote(request: Request) -> Redis:
return request.app.state.redis_remote
return request.app.state.redis_market
def get_redis(request: Request) -> Redis:
return request.app.state.redis_remote
redis_cache = getattr(request.app.state,'redis_cache',None)
if redis_cache is None:
raise RuntimeError('Cache KeyDB is not configured')
return redis_cache
def get_cache_remote(redis_client: Redis = Depends(get_redis_remote)) -> ICache:
return KeydbCache(redis_client)
def get_cache_remote(request: Request) -> ICache:
return KeydbCache(get_redis(request))
def get_remote_cache(redis_client: Redis = Depends(get_redis_remote)) -> ICache:
return RemoteCache(redis_client)
def get_remote_cache(request: Request) -> ICache:
return RemoteCache(get_redis_remote(request))
def get_cache(cache: ICache = Depends(get_cache_remote)) -> ICache: