Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions docs/modules/background_tasks.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`.
12 changes: 12 additions & 0 deletions docs/modules/settings.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.<package>.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.<package>.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
Expand Down
2 changes: 2 additions & 0 deletions modules/background_tasks/background_tasks/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
]
57 changes: 57 additions & 0 deletions modules/background_tasks/background_tasks/worker_settings.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
"""Hydrated DB-backed module settings for task code (GH #378).

A Celery worker never runs the hosting lifespan that hydrates
``app.state.<package>.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:
# 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.model_copy(deep=True)


def clear_settings_cache() -> None:
"""Drop every cached settings object (tests, or after a known write)."""
with _lock:
_cache.clear()
56 changes: 56 additions & 0 deletions modules/background_tasks/tests/test_bg_worker_settings.py
Original file line number Diff line number Diff line change
@@ -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"
27 changes: 23 additions & 4 deletions modules/settings/settings/hydrate.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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))
52 changes: 42 additions & 10 deletions modules/settings/settings/store.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,15 +13,56 @@

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


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 ``<package>.<field>`` 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(session.execute(stmt).tuples(), package)


class SettingsStore:
"""SYSTEM-scoped key/value store keyed by ``(package, field)``."""

Expand All @@ -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.
Expand Down
26 changes: 26 additions & 0 deletions modules/settings/tests/test_hydrate.py
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Loading