from typing import Any, Literal from fansifter_common.utils.functional import lazy_proxy from pydantic import Field, SecretStr, ValidationError from pydantic_settings import BaseSettings from sqlalchemy import URL from app.dsp.enums import DSPClientName class Settings(BaseSettings): model_config = { "env_file": ".env", "extra": "ignore", } # ── App ────────────────────────────────────────────────────────────── environment: str = Field("dev", validation_alias="ENVIRONMENT") debug: bool = Field(False, validation_alias="APP_DEBUG") cors_origins: list[str] = Field( ["http://localhost:5173"], validation_alias="CORS_ORIGINS" ) service_name: str = "resonance-engine" service_version: str = "1.0.0" # ── AWS ────────────────────────────────────────────────────────────── aws_region_name: str = Field("us-east-1", validation_alias="AWS_REGION_NAME") # ── KMS / Encrypter ────────────────────────────────────────────────── kms_key_arn: str = Field("", validation_alias="KMS_KEY_ARN") encrypter_backend: Literal["fernet", "kms"] = Field( "kms", validation_alias="ENCRYPTER_BACKEND" ) fernet_keys: list[str] = Field(default_factory=list, validation_alias="FERNET_KEYS") fernet_keys_secret_name: str = Field("", validation_alias="FERNET_KEYS_SECRET_NAME") # ── Aurora DSQL ────────────────────────────────────────────────────── dsql_endpoint: str | None = Field(None, validation_alias="DSQL_ENDPOINT") dsql_token_expires_in: int = Field(900, validation_alias="DSQL_TOKEN_EXPIRES_IN") # ── Database ───────────────────────────────────────────────────────── db_user: str = Field("resonance", validation_alias="DB_USER") db_password: SecretStr | None = Field(None, validation_alias="DB_PASSWORD") db_host: str = Field("localhost", validation_alias="DB_HOST") db_port: int = Field(5432, validation_alias="DB_PORT") db_name: str = Field("resonance", validation_alias="DB_NAME") db_echo: bool = Field(False, validation_alias="DB_ECHO") db_pool_size: int = Field(5, validation_alias="DB_POOL_SIZE") db_pool_max_overflow: int = Field(10, validation_alias="DB_POOL_MAX_OVERFLOW") db_pool_recycle: int = Field(300, validation_alias="DB_POOL_RECYCLE") @property def db_url(self) -> URL: return URL.create( drivername="postgresql+psycopg", username=self.db_user, password=self.db_password.get_secret_value() if self.db_password else None, host=self.db_host, port=self.db_port, database=self.db_name, ) # ── Redis ───────────────────────────────────────────────────────────── redis_url: str = Field("redis://localhost:6379/0", validation_alias="REDIS_URL") # ── Kafka ───────────────────────────────────────────────────────────── kafka_bootstrap_servers: str = Field("", validation_alias="KAFKA_BOOTSTRAP_SERVERS") kafka_security_protocol: str = Field( "SSL", validation_alias="KAFKA_SECURITY_PROTOCOL" ) kafka_topic_fans: str = Field( "resonance-engine.fan-profile", validation_alias="KAFKA_TOPIC_FANS" ) kafka_topic_top_artists: str = Field( "resonance-engine.fan-top-artists", validation_alias="KAFKA_TOPIC_TOP_ARTISTS" ) kafka_topic_top_tracks: str = Field( "resonance-engine.fan-top-tracks", validation_alias="KAFKA_TOPIC_TOP_TRACKS" ) kafka_topic_recently_played: str = Field( "resonance-engine.fan-recently-played", validation_alias="KAFKA_TOPIC_RECENTLY_PLAYED", ) kafka_topic_playlists: str = Field( "resonance-engine.fan-playlists", validation_alias="KAFKA_TOPIC_PLAYLISTS" ) kafka_topic_saved_albums: str = Field( "resonance-engine.fan-saved-albums", validation_alias="KAFKA_TOPIC_SAVED_ALBUMS" ) kafka_topic_saved_tracks: str = Field( "resonance-engine.fan-saved-tracks", validation_alias="KAFKA_TOPIC_SAVED_TRACKS" ) kafka_topic_followed_artists: str = Field( "resonance-engine.fan-followed-artists", validation_alias="KAFKA_TOPIC_FOLLOWED_ARTISTS", ) # ── Spotify — Songwhip ─────────────────────────────────────────────── spotify_songwhip_client_id: str = Field( "", validation_alias="SPOTIFY_SONGWHIP_CLIENT_ID" ) spotify_songwhip_client_secret: str = Field( "", validation_alias="SPOTIFY_SONGWHIP_CLIENT_SECRET" ) spotify_songwhip_use_faker: bool = Field( True, validation_alias="SPOTIFY_SONGWHIP_USE_FAKER" ) # ── Spotify — SMF ──────────────────────────────────────────────────── spotify_smf_client_id: str = Field("", validation_alias="SPOTIFY_SMF_CLIENT_ID") spotify_smf_client_secret: str = Field( "", validation_alias="SPOTIFY_SMF_CLIENT_SECRET" ) spotify_smf_use_faker: bool = Field(True, validation_alias="SPOTIFY_SMF_USE_FAKER") # ── DSP ────────────────────────────────────────────────────────────── dsp_rate_limit_stats_backend: Literal["redis", "memory"] = Field( "memory", validation_alias="DSP_RATE_LIMIT_STATS_BACKEND" ) # ── Data sink ──────────────────────────────────────────────────────── data_sink_backend: Literal["postgres", "kafka", "dummy"] = Field( "postgres", validation_alias="DATA_SINK_BACKEND" ) # ── Fan fanout ─────────────────────────────────────────────────────── fan_fanout_batch_size: int = Field(100, validation_alias="FAN_FANOUT_BATCH_SIZE") fan_fanout_spotify_smf_batch_size: int | None = Field( None, validation_alias="FAN_FANOUT_SPOTIFY_SMF_BATCH_SIZE" ) fan_fanout_spotify_songwhip_batch_size: int | None = Field( None, validation_alias="FAN_FANOUT_SPOTIFY_SONGWHIP_BATCH_SIZE" ) fan_fanout_time_limit_s: int = Field( 900, validation_alias="FAN_FANOUT_TIME_LIMIT_S" ) fan_fanout_max_fans: int = Field(20_000_000, validation_alias="FAN_FANOUT_MAX_FANS") fan_fanout_dispatch_backend: Literal["lambda", "dummy"] = Field( "lambda", validation_alias="FAN_FANOUT_DISPATCH_BACKEND" ) fan_fanout_lambda_function_name: str = Field( "", validation_alias="FAN_FANOUT_LAMBDA_FUNCTION_NAME" ) def fan_fanout_batch_size_for(self, client_name: DSPClientName) -> int: if ( client_name == DSPClientName.spotify_smf and self.fan_fanout_spotify_smf_batch_size is not None ): return self.fan_fanout_spotify_smf_batch_size if ( client_name == DSPClientName.spotify_songwhip and self.fan_fanout_spotify_songwhip_batch_size is not None ): return self.fan_fanout_spotify_songwhip_batch_size return self.fan_fanout_batch_size # ── Fan collect ────────────────────────────────────────────────────── fan_collect_dispatch_backend: Literal["sqs", "dummy"] = Field( "sqs", validation_alias="FAN_COLLECT_DISPATCH_BACKEND" ) fan_collect_time_limit_s: int = Field( 600, validation_alias="FAN_COLLECT_TIME_LIMIT_S" ) fan_collect_timeout_buffer_s: int = Field( 30, validation_alias="FAN_COLLECT_TIMEOUT_BUFFER_S" ) fan_collect_rate_limit_circuit_breaker_threshold: int = Field( 3, validation_alias="FAN_COLLECT_RATE_LIMIT_CIRCUIT_BREAKER_THRESHOLD" ) fan_collect_concurrency: int = Field(15, validation_alias="FAN_COLLECT_CONCURRENCY") fan_collect_spotify_smf_concurrency: int = Field( 15, validation_alias="FAN_COLLECT_SPOTIFY_SMF_CONCURRENCY" ) fan_collect_spotify_songwhip_concurrency: int = Field( 5, validation_alias="FAN_COLLECT_SPOTIFY_SONGWHIP_CONCURRENCY" ) # Fan collect — per-step re-collect intervals fan_collect_profile_interval_s: int = Field( 0, validation_alias="FAN_COLLECT_PROFILE_INTERVAL_S" ) fan_collect_top_artists_interval_s: int = Field( 0, validation_alias="FAN_COLLECT_TOP_ARTISTS_INTERVAL_S" ) fan_collect_top_tracks_interval_s: int = Field( 0, validation_alias="FAN_COLLECT_TOP_TRACKS_INTERVAL_S" ) fan_collect_recently_played_interval_s: int = Field( 0, validation_alias="FAN_COLLECT_RECENTLY_PLAYED_INTERVAL_S" ) fan_collect_playlists_interval_s: int = Field( 0, validation_alias="FAN_COLLECT_PLAYLISTS_INTERVAL_S" ) fan_collect_saved_albums_interval_s: int = Field( 0, validation_alias="FAN_COLLECT_SAVED_ALBUMS_INTERVAL_S" ) fan_collect_saved_tracks_interval_s: int = Field( 0, validation_alias="FAN_COLLECT_SAVED_TRACKS_INTERVAL_S" ) fan_collect_followed_artists_interval_s: int = Field( 0, validation_alias="FAN_COLLECT_FOLLOWED_ARTISTS_INTERVAL_S" ) @property def fan_collect_cooldown_s(self) -> int: """Minimum of all step intervals; 0 when any interval is 0 (cooldown disabled).""" intervals = [ self.fan_collect_profile_interval_s, self.fan_collect_top_artists_interval_s, self.fan_collect_top_tracks_interval_s, self.fan_collect_recently_played_interval_s, self.fan_collect_playlists_interval_s, self.fan_collect_saved_albums_interval_s, self.fan_collect_saved_tracks_interval_s, self.fan_collect_followed_artists_interval_s, ] if any(i == 0 for i in intervals): return 0 return min(intervals) def fan_collect_concurrency_for(self, client_name: DSPClientName) -> int: if client_name == DSPClientName.spotify_songwhip: return self.fan_collect_spotify_songwhip_concurrency if client_name == DSPClientName.spotify_smf: return self.fan_collect_spotify_smf_concurrency return self.fan_collect_concurrency # ── SQS queues ─────────────────────────────────────────────────────── sqs_fan_collect_spotify_smf_queue_url: str = Field( "", validation_alias="SQS_FAN_COLLECT_SPOTIFY_SMF_QUEUE_URL" ) sqs_fan_collect_spotify_songwhip_queue_url: str = Field( "", validation_alias="SQS_FAN_COLLECT_SPOTIFY_SONGWHIP_QUEUE_URL" ) def sqs_fan_collect_queue_url_for(self, client_name: DSPClientName) -> str: if client_name == DSPClientName.spotify_smf: return self.sqs_fan_collect_spotify_smf_queue_url if client_name == DSPClientName.spotify_songwhip: return self.sqs_fan_collect_spotify_songwhip_queue_url raise ValueError(f"Unknown DSP client: {client_name}") # ── Celery queues ──────────────────────────────────────────────────── celery_fan_fanout_queue: str = Field( "resonance.fan_fanout", validation_alias="CELERY_FAN_FANOUT_QUEUE" ) celery_fan_collect_spotify_smf_queue: str = Field( "resonance.fan_collect.spotify_smf", validation_alias="CELERY_FAN_COLLECT_SPOTIFY_SMF_QUEUE", ) celery_fan_collect_spotify_songwhip_queue: str = Field( "resonance.fan_collect.spotify_songwhip", validation_alias="CELERY_FAN_COLLECT_SPOTIFY_SONGWHIP_QUEUE", ) def celery_fan_collect_queue_for(self, client_name: DSPClientName) -> str: if client_name == DSPClientName.spotify_smf: return self.celery_fan_collect_spotify_smf_queue if client_name == DSPClientName.spotify_songwhip: return self.celery_fan_collect_spotify_songwhip_queue raise ValueError(f"Unknown DSP client: {client_name}") # ── Pipeline run ───────────────────────────────────────────────────── pipeline_run_stale_timeout_s: int = Field( 1800, validation_alias="PIPELINE_RUN_STALE_TIMEOUT_S" ) pipeline_run_estimate_avg_latency_ms: int = Field( 135, validation_alias="PIPELINE_RUN_ESTIMATE_AVG_LATENCY_MS" ) # ── Logging ────────────────────────────────────────────────────────── logging_debug: bool = Field(False, validation_alias="LOGGING_DEBUG") @property def logging_config(self) -> dict[str, Any]: return { "version": 1, "disable_existing_loggers": False, "formatters": { "json": { "()": "owslogger.logger.DDJsonFormatter", "service_name": self.service_name, "service_version": self.service_version, "env": self.environment, }, "console": { "()": "fansifter_common.logging.ConsoleFormatter", }, }, "filters": { "require_debug_true": { "()": "fansifter_common.logging.RequireDebugTrueFilter", "value": self.logging_debug, }, "require_debug_false": { "()": "fansifter_common.logging.RequireDebugFalseFilter", "value": self.logging_debug, }, }, "handlers": { "stream": { "level": "DEBUG", "class": "logging.StreamHandler", "formatter": "console", "filters": ["require_debug_false"], }, "rich": { "level": "DEBUG", "class": "rich.logging.RichHandler", "filters": ["require_debug_true"], "formatter": "console", }, "null": { "level": "DEBUG", "class": "logging.NullHandler", }, }, "loggers": { "app": { "handlers": ["stream", "rich"], "level": "INFO", }, "sqlalchemy.engine": { "handlers": ["stream", "rich"], "level": "WARNING", }, "httpx": { "handlers": ["stream", "rich"], "level": "WARNING", }, "uvicorn": { "handlers": ["stream", "rich"], "level": "INFO", }, "uvicorn.access": { "handlers": ["null"], "level": "INFO", }, "celery": { "handlers": ["stream", "rich"], "level": "INFO", }, "celery.worker.strategy": { "handlers": ["stream", "rich"], "level": "WARNING", "propagate": False, }, "celery.app.trace": { "handlers": ["stream", "rich"], "level": "WARNING", "propagate": False, }, }, } def get_settings(**defaults: Any) -> Settings: """Setup the application settings.""" try: return Settings(**defaults) except ValidationError as exc: errors = "\n".join( [ f" * {'.'.join(map(str, error['loc']))} - {error['msg']}" for error in exc.errors() ] ) raise RuntimeError(f"Failed to initialize settings:\n {errors}") from exc settings = lazy_proxy(get_settings)