feat: change keydb providing
This commit is contained in:
@@ -1,7 +0,0 @@
|
||||
BEGIN;
|
||||
|
||||
ALTER TABLE sbp_withdrawals
|
||||
DROP COLUMN IF EXISTS wallet_address,
|
||||
DROP COLUMN IF EXISTS sender_wallet_address;
|
||||
|
||||
COMMIT;
|
||||
@@ -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:
|
||||
|
||||
2
src/infrastructure/cache/client.py
vendored
2
src/infrastructure/cache/client.py
vendored
@@ -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,
|
||||
|
||||
@@ -1,14 +1,14 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from functools import lru_cache
|
||||
from typing import List, Literal
|
||||
import os
|
||||
from dotenv import load_dotenv, find_dotenv
|
||||
from pydantic import AliasChoices, Field, field_validator, model_validator
|
||||
from pydantic_settings import BaseSettings, SettingsConfigDict
|
||||
from functools import lru_cache
|
||||
from typing import List,Literal
|
||||
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:
|
||||
|
||||
34
src/main.py
34
src/main.py
@@ -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}')
|
||||
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user