From f4daa4c6b0cffbb44640c390a1e3da822d069d86 Mon Sep 17 00:00:00 2001 From: Anto Subash Date: Mon, 5 Oct 2026 12:33:20 +0200 Subject: [PATCH 1/2] fix(settings): sync hydration of DB-backed settings for Celery workers (#378) Add settings.hydrate.hydrate_settings_sync and background_tasks.settings_for, sharing key mapping and value_type parsing with the async path. Claude-Session: https://claude.ai/code/session_01M9neheZZEe3sVpDi2S3zT4 --- docs/modules/background_tasks.md | 16 ++++++ docs/modules/settings.md | 12 ++++ .../background_tasks/__init__.py | 2 + .../background_tasks/worker_settings.py | 55 ++++++++++++++++++ .../tests/test_bg_worker_settings.py | 56 +++++++++++++++++++ modules/settings/settings/hydrate.py | 27 +++++++-- modules/settings/settings/store.py | 52 +++++++++++++---- modules/settings/tests/test_hydrate.py | 26 +++++++++ 8 files changed, 232 insertions(+), 14 deletions(-) create mode 100644 modules/background_tasks/background_tasks/worker_settings.py create mode 100644 modules/background_tasks/tests/test_bg_worker_settings.py diff --git a/docs/modules/background_tasks.md b/docs/modules/background_tasks.md index 29c10a9a..9af5d253 100644 --- a/docs/modules/background_tasks.md +++ b/docs/modules/background_tasks.md @@ -220,3 +220,19 @@ Production deployments running multiple workers can write concurrently because e ## Locales Top-level keys in `background_tasks/locales/en.json`: `index`, `filters`, `status` (per-state labels), `table`, `detail`, `retry_dialog`, `toasts`. + +## Reading DB-backed settings in a task + +Workers never run the hosting lifespan, so `FileStorageSettings()` (or any `DbBackedSettings`) built inside a task returns pydantic defaults, not what the admin UI saved. Use `background_tasks.settings_for(cls, package)`: + +```python +from background_tasks import settings_for +from file_storage.settings import FileStorageSettings + + +@celery.task +def convert(file_id: int) -> None: + cfg = settings_for(FileStorageSettings, "file_storage") +``` + +It reads SYSTEM-scope overrides through the worker's sync session and caches per process for 30 seconds (`ttl=0` always re-reads; `clear_settings_cache()` drops the cache). The underlying sync helper is `settings.hydrate.hydrate_settings_sync`. diff --git a/docs/modules/settings.md b/docs/modules/settings.md index 95c765fd..2759fb52 100644 --- a/docs/modules/settings.md +++ b/docs/modules/settings.md @@ -27,6 +27,18 @@ The pattern (pydantic `BaseSettings` subclass + `register_module_settings` in `r - **Admin editing**: registered fields appear under that package at `/admin/settings/` with type-aware inputs. - **Hot reload**: saving via the admin UI calls `apply_changes_and_reload`, which validates the diff against the pydantic class, persists deltas, swaps the live `app.state..settings`, and publishes [`SettingsReloaded`](#events) so dependents (SMTP clients, Celery configs, …) can rebuild. +### Read DB-backed settings outside the web process + +`DbBackedSettings` subclasses ignore the environment, and only the hosting lifespan hydrates `app.state..settings`. A Celery worker (or CLI script) never runs that lifespan, so a bare `FileStorageSettings()` there silently returns pydantic defaults. Hydrate synchronously instead: + +```python +from settings.hydrate import hydrate_settings_sync + +cfg = hydrate_settings_sync(FileStorageSettings, session, "file_storage") +``` + +`session` is a plain sync SQLAlchemy `Session`. Precedence matches the web process (SYSTEM-scope DB overrides over pydantic defaults, same `value_type` parsing, no env). Inside a Celery task prefer `background_tasks.settings_for(FileStorageSettings, "file_storage")`, which opens the worker's sync session and caches the result per process for 30 seconds (`ttl=` to change, `0` to always re-read). + ### Read settings at request time (generic K/V) ```python diff --git a/modules/background_tasks/background_tasks/__init__.py b/modules/background_tasks/background_tasks/__init__.py index f1a81e72..e7e2e9dd 100644 --- a/modules/background_tasks/background_tasks/__init__.py +++ b/modules/background_tasks/background_tasks/__init__.py @@ -5,9 +5,11 @@ get_log_context, install_log_filter, ) +from background_tasks.worker_settings import settings_for __all__ = [ "bind_task_context", "get_log_context", "install_log_filter", + "settings_for", ] diff --git a/modules/background_tasks/background_tasks/worker_settings.py b/modules/background_tasks/background_tasks/worker_settings.py new file mode 100644 index 00000000..1d22ac31 --- /dev/null +++ b/modules/background_tasks/background_tasks/worker_settings.py @@ -0,0 +1,55 @@ +"""Hydrated DB-backed module settings for task code (GH #378). + +A Celery worker never runs the hosting lifespan that hydrates +``app.state..settings``, and ``DbBackedSettings`` ignores the +environment, so a bare ``FileStorageSettings()`` in a task returns pydantic +defaults without complaint. Use :func:`settings_for` instead:: + + cfg = settings_for(FileStorageSettings, "file_storage") + +It reads the SYSTEM-scope overrides through the worker's sync session and +caches the result per process for ``DEFAULT_TTL_SECONDS`` so a hot task does not +query on every call. Admin edits therefore reach a worker within the TTL. +""" + +from __future__ import annotations + +import threading +import time +from typing import cast + +from pydantic_settings import BaseSettings +from settings.hydrate import hydrate_settings_sync + +from background_tasks.sync_db import sync_session + +DEFAULT_TTL_SECONDS = 30.0 + +_lock = threading.Lock() +_cache: dict[tuple[type, str], tuple[float, BaseSettings]] = {} + + +def settings_for[T: BaseSettings]( + cls: type[T], package: str, *, ttl: float = DEFAULT_TTL_SECONDS +) -> T: + """Return ``cls`` hydrated from DB overrides, cached per process for ``ttl`` seconds. + + ``ttl=0`` always re-reads. + """ + key = (cls, package) + now = time.monotonic() + with _lock: + hit = _cache.get(key) + if hit is not None and ttl > 0 and now - hit[0] < ttl: + return cast(T, hit[1]) + with sync_session() as session: + value = hydrate_settings_sync(cls, session, package) + with _lock: + _cache[key] = (now, value) + return value + + +def clear_settings_cache() -> None: + """Drop every cached settings object (tests, or after a known write).""" + with _lock: + _cache.clear() diff --git a/modules/background_tasks/tests/test_bg_worker_settings.py b/modules/background_tasks/tests/test_bg_worker_settings.py new file mode 100644 index 00000000..1c41bf89 --- /dev/null +++ b/modules/background_tasks/tests/test_bg_worker_settings.py @@ -0,0 +1,56 @@ +"""settings_for hydrates DB-backed settings through the worker's sync session.""" + +from __future__ import annotations + +from collections.abc import Iterator + +import pytest +from background_tasks import sync_db +from background_tasks.worker_settings import clear_settings_cache, settings_for +from pydantic import Field +from settings.models import Setting +from simple_module_core.settings_base import DbBackedSettings +from sqlalchemy import create_engine +from sqlalchemy.orm import Session + + +class _Cfg(DbBackedSettings): + backend: str = "filesystem" + tags: list[str] = Field(default_factory=list) + + +@pytest.fixture +def db_url(tmp_path) -> Iterator[str]: + url = f"sqlite:///{tmp_path / 'w.db'}" + engine = create_engine(url) + Setting.metadata.create_all(engine) + sync_db.set_database_url(url) + clear_settings_cache() + yield url + sync_db.dispose_sync_engine() + clear_settings_cache() + engine.dispose() + + +def _put(url: str, key: str, value: str, vtype: str) -> None: + engine = create_engine(url) + with Session(engine) as s: + s.add(Setting(key=key, value=value, value_type=vtype)) + s.commit() + engine.dispose() + + +def test_settings_for_reads_db_overrides(db_url: str) -> None: + assert _Cfg().backend == "filesystem" + _put(db_url, "demo.backend", "s3", "string") + _put(db_url, "demo.tags", '["a","b"]', "json") + cfg = settings_for(_Cfg, "demo") + assert cfg.backend == "s3" + assert cfg.tags == ["a", "b"] + + +def test_settings_for_caches_until_ttl_expires(db_url: str) -> None: + assert settings_for(_Cfg, "demo").backend == "filesystem" + _put(db_url, "demo.backend", "s3", "string") + assert settings_for(_Cfg, "demo").backend == "filesystem" # cached + assert settings_for(_Cfg, "demo", ttl=0).backend == "s3" diff --git a/modules/settings/settings/hydrate.py b/modules/settings/settings/hydrate.py index 14fff1b6..c81a9e86 100644 --- a/modules/settings/settings/hydrate.py +++ b/modules/settings/settings/hydrate.py @@ -13,8 +13,9 @@ from typing import get_origin from pydantic_settings import BaseSettings +from sqlalchemy.orm import Session -from settings.store import SettingsStore +from settings.store import SettingsStore, get_overrides_sync def value_type_for_field(cls: type[BaseSettings], field_name: str) -> str: @@ -52,12 +53,30 @@ def _parse(raw: str, value_type: str): return raw -async def hydrate_settings[T: BaseSettings](cls: type[T], store: SettingsStore, package: str) -> T: - """Construct ``cls`` with DB overrides merged over pydantic defaults.""" - raw_overrides = await store.get_overrides(package) +def build_settings[T: BaseSettings](cls: type[T], raw_overrides: dict[str, tuple[str, str]]) -> T: + """Construct ``cls`` from ``{field: (raw, value_type)}``, skipping unknown fields. + + Shared by the async and sync hydrators so parsing cannot drift. + """ parsed: dict[str, object] = {} for field_name, (raw, vtype) in raw_overrides.items(): if field_name not in cls.model_fields: continue parsed[field_name] = _parse(raw, vtype) return cls(**parsed) + + +async def hydrate_settings[T: BaseSettings](cls: type[T], store: SettingsStore, package: str) -> T: + """Construct ``cls`` with DB overrides merged over pydantic defaults.""" + return build_settings(cls, await store.get_overrides(package)) + + +def hydrate_settings_sync[T: BaseSettings](cls: type[T], session: Session, package: str) -> T: + """Sync :func:`hydrate_settings` for Celery workers and other non-async code. + + Workers never run the hosting lifespan, so a bare ``cls()`` there returns + pydantic defaults (``DbBackedSettings`` ignores the environment on purpose). + Takes a plain sync ``Session`` and applies the same SYSTEM-scope overrides + with the same ``value_type`` parsing as the web process. + """ + return build_settings(cls, get_overrides_sync(session, package)) diff --git a/modules/settings/settings/store.py b/modules/settings/settings/store.py index 7618ff40..371aa155 100644 --- a/modules/settings/settings/store.py +++ b/modules/settings/settings/store.py @@ -13,8 +13,14 @@ from __future__ import annotations +from collections.abc import Iterable + +from sqlalchemy import select +from sqlalchemy.orm import Session + from settings.constants import SYSTEM_SCOPE_ID from settings.contracts.schemas import SettingScope, SettingUpsert, SettingValueType +from settings.models import Setting from settings.service import SettingService @@ -22,6 +28,41 @@ def _key(package: str, field: str) -> str: return f"{package}.{field}" +def package_overrides( + rows: Iterable[tuple[str, str, str]], package: str +) -> dict[str, tuple[str, str]]: + """Map ``(key, value, value_type)`` rows to ``{field_name: (raw, value_type)}``. + + The one place the ``.`` key format is interpreted, shared by + the async :meth:`SettingsStore.get_overrides` and the sync + :func:`get_overrides_sync` so the two cannot drift. + """ + prefix = f"{package}." + out: dict[str, tuple[str, str]] = {} + for key, value, value_type in rows: + if not key.startswith(prefix): + continue + field_name = key[len(prefix) :] + if "." in field_name: + continue + out[field_name] = (value, value_type) + return out + + +def get_overrides_sync(session: Session, package: str) -> dict[str, tuple[str, str]]: + """Sync twin of :meth:`SettingsStore.get_overrides` for worker processes. + + Reads SYSTEM-scope rows through a plain sync ``Session`` (unmasked, like the + async path — this feeds live settings objects, not a screen). + """ + stmt = select(Setting.key, Setting.value, Setting.value_type).where( + Setting.scope == SettingScope.SYSTEM.value, + Setting.scope_id == SYSTEM_SCOPE_ID, + Setting.key.startswith(f"{package}.", autoescape=True), + ) + return package_overrides(((k, v, t) for k, v, t in session.execute(stmt)), package) + + class SettingsStore: """SYSTEM-scoped key/value store keyed by ``(package, field)``.""" @@ -30,17 +71,8 @@ def __init__(self, service: SettingService) -> None: async def get_overrides(self, package: str) -> dict[str, tuple[str, str]]: """Return ``{field_name: (raw_value, value_type)}`` for a package.""" - prefix = f"{package}." items = await self._service.list_by_scope_unmasked(SettingScope.SYSTEM, SYSTEM_SCOPE_ID) - out: dict[str, tuple[str, str]] = {} - for item in items: - if not item.key.startswith(prefix): - continue - field_name = item.key[len(prefix) :] - if "." in field_name: - continue - out[field_name] = (item.value, item.value_type) - return out + return package_overrides(((i.key, i.value, i.value_type) for i in items), package) async def all_override_fields(self) -> dict[str, frozenset[str]]: """Return ``{package: {field_name, ...}}`` for every stored override. diff --git a/modules/settings/tests/test_hydrate.py b/modules/settings/tests/test_hydrate.py index 6630cb1b..eaf5e26f 100644 --- a/modules/settings/tests/test_hydrate.py +++ b/modules/settings/tests/test_hydrate.py @@ -63,3 +63,29 @@ def test_value_type_for_bool_int_float_str_list() -> None: assert value_type_for_field(_Cfg, "rate") == "float" assert value_type_for_field(_Cfg, "host") == "string" assert value_type_for_field(_Cfg, "tags") == "json" + + +def test_hydrate_sync_matches_async_parsing(tmp_path) -> None: + from settings.hydrate import hydrate_settings_sync + from settings.models import Setting + from sqlalchemy import create_engine + from sqlalchemy.orm import Session + + engine = create_engine(f"sqlite:///{tmp_path / 's.db'}") + Setting.metadata.create_all(engine) + with Session(engine) as session: + for k, v, t in [ + ("demo.port", "587", "int"), + ("demo.allow", "true", "bool"), + ("demo.tags", '["x","y"]', "json"), + ("demo.unknown", "1", "int"), + ("demo.a.b", "1", "int"), + ("other.port", "1", "int"), + ]: + session.add(Setting(key=k, value=v, value_type=t)) + session.add( + Setting(scope="user", scope_id="7", key="demo.host", value="no", value_type="string") + ) + session.commit() + cfg = hydrate_settings_sync(_Cfg, session, "demo") + assert (cfg.port, cfg.allow, cfg.tags, cfg.host) == (587, True, ["x", "y"], "localhost") From ecb26135753babcc283bc2696993ee3d7ef29794 Mon Sep 17 00:00:00 2001 From: Anto Subash Date: Tue, 6 Oct 2026 23:18:31 +0200 Subject: [PATCH 2/2] fix: address code review findings (pass 1) Claude-Session: https://claude.ai/code/session_01M9neheZZEe3sVpDi2S3zT4 --- .../background_tasks/background_tasks/worker_settings.py | 6 ++++-- modules/settings/settings/store.py | 2 +- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/modules/background_tasks/background_tasks/worker_settings.py b/modules/background_tasks/background_tasks/worker_settings.py index 1d22ac31..b94e2474 100644 --- a/modules/background_tasks/background_tasks/worker_settings.py +++ b/modules/background_tasks/background_tasks/worker_settings.py @@ -41,12 +41,14 @@ def settings_for[T: BaseSettings]( with _lock: hit = _cache.get(key) if hit is not None and ttl > 0 and now - hit[0] < ttl: - return cast(T, hit[1]) + # A copy: BaseSettings is mutable, so one task editing its settings + # must not change what the next task in this process sees. + return cast(T, hit[1]).model_copy(deep=True) with sync_session() as session: value = hydrate_settings_sync(cls, session, package) with _lock: _cache[key] = (now, value) - return value + return value.model_copy(deep=True) def clear_settings_cache() -> None: diff --git a/modules/settings/settings/store.py b/modules/settings/settings/store.py index 371aa155..619e5adb 100644 --- a/modules/settings/settings/store.py +++ b/modules/settings/settings/store.py @@ -60,7 +60,7 @@ def get_overrides_sync(session: Session, package: str) -> dict[str, tuple[str, s Setting.scope_id == SYSTEM_SCOPE_ID, Setting.key.startswith(f"{package}.", autoescape=True), ) - return package_overrides(((k, v, t) for k, v, t in session.execute(stmt)), package) + return package_overrides(session.execute(stmt).tuples(), package) class SettingsStore: