diff --git a/CHANGELOG.md b/CHANGELOG.md index 227ceb8..93c28c8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -39,6 +39,19 @@ This project adheres to [Semantic Versioning](https://semver.org/). --> +## [1.2.3] = 2026-10-07 +### Changed +- ⚡ updated database drivers to use psycopg2 in server +- ⚡ pinned geoalchemy2 to `>=0.18.0,<1` to comply with sqlalchemy requirements + +### Added +- 🔥 docker bake definition file for migrations and backend +- 🔥 value backfill to migrate numeric `value_jsonb` to `value_float` for existing data in `probe_data` table + +### Fixed +- 🩹 Bug in load data related to metric_type +- 🩹 Bug in server/backend for loading api keys + ## [1.2.2] = 2026-09-30 ### Changed - ⚡ `probe_data` table now has `value_float` and `value_jsonb` columns instead of one `value` (jsonb type) column. diff --git a/docs/guides/index.md b/docs/guides/index.md index d2a237f..7265702 100644 --- a/docs/guides/index.md +++ b/docs/guides/index.md @@ -3,6 +3,7 @@ * [Configuration](configuration.md) * [Using the `opensampl` CLI](opensampl-cli.md) * [Using the `opensampl-server` CLI](opensampl-server.md) +* [Value Backfill](value-backfill.md) * [Collection Guide](collection.md) * [Random Data Guide](random-data-generation.md) * [NTP Extension Guide](ntp-extension.md) diff --git a/docs/guides/opensampl-cli.md b/docs/guides/opensampl-cli.md index cf27934..a705e3c 100644 --- a/docs/guides/opensampl-cli.md +++ b/docs/guides/opensampl-cli.md @@ -3,6 +3,20 @@ Use `opensampl --help` to see the top-level commands, or `opensampl --help` for subcommand-specific options. +## Maintenance + +Database maintenance commands are grouped under `opensampl maintenance`. The +`value-backfill` command copies historical numeric JSON values into the +optimized floating-point column introduced by the probe-data type migration: + +```bash +opensampl maintenance value-backfill +``` + +It requires a direct database connection and refuses to run when +`ROUTE_TO_BACKEND=true`. See the [value backfill guide](value-backfill.md) for +prerequisites, tuning options, Compose usage, and recovery instructions. + ## Load Data ### Probe Data @@ -114,4 +128,3 @@ Arguments: Options: * `--update-db` (`-u`): Update the database with the new probe type - diff --git a/docs/guides/opensampl-server.md b/docs/guides/opensampl-server.md index 5755e0c..1906028 100644 --- a/docs/guides/opensampl-server.md +++ b/docs/guides/opensampl-server.md @@ -78,6 +78,17 @@ opensampl-server run backend python -m opensampl.cli init This maps directly to `docker compose run --rm ...`. +The packaged deployment also contains an opt-in service for the historical +numeric value backfill: + +```bash +opensampl-server run value-backfill +``` + +The service is not started by the normal `up` command. See the +[value backfill guide](value-backfill.md) before running it, especially for a +large database. + ## Using a custom env file `--env-file` is a top-level CLI option, so it must appear before the subcommand: diff --git a/docs/guides/value-backfill.md b/docs/guides/value-backfill.md new file mode 100644 index 0000000..5cd7b5a --- /dev/null +++ b/docs/guides/value-backfill.md @@ -0,0 +1,125 @@ +# Value backfill + +Migration `b88042bae240` changes how OpenSAMPL stores numeric probe values in +`castdb.probe_data`. It renames the original `value` column to `value_jsonb` and +adds an optimized `value_float` column. The migration only changes the schema; +it intentionally does not rewrite historical data because that can take much +longer than a normal deployment migration. + +The value backfill performs that historical rewrite as a separate maintenance +operation. It: + +- finds metric types whose `value_type` is `float` or `int`; +- converts their scalar `value_jsonb` values to double precision and stores the + result in `value_float`; +- processes the table in newest-first time windows using independent + transactions and configurable parallel workers; and +- vacuums `castdb.probe_data` periodically during the backfill. + +Only rows where `value_float IS NULL` and `value_jsonb IS NOT NULL` are updated. +This makes the value backfill resumable. If it stops or fails, correct the +problem and run the same command again; previously backfilled rows are skipped. + +## Before running + +Run the value backfill only after the deployment's Alembic migrations have +completed. The command verifies that both `value_jsonb` and `value_float` exist +before starting. + +The value backfill requires a direct database connection through +`DATABASE_URL`. It does not use the OpenSAMPL backend API. If +`ROUTE_TO_BACKEND=true`, it prints a warning and exits without connecting. +Explicitly set `ROUTE_TO_BACKEND=false` for the maintenance run. + +Updates and vacuum operations can generate substantial database I/O. For a +large deployment: + +- verify that a recent backup is available; +- run during a maintenance or low-traffic period; +- begin with a conservative worker count; and +- do not run multiple value backfills concurrently. + +## Run through the OpenSAMPL CLI + +With `DATABASE_URL` configured and `ROUTE_TO_BACKEND=false`, run: + +```bash +opensampl maintenance value-backfill +``` + +To select a particular OpenSAMPL environment file: + +```bash +opensampl --env-file ./maintenance.env maintenance value-backfill +``` + +The required settings can also be supplied for a one-off shell invocation: + +```bash +ROUTE_TO_BACKEND=false \ +DATABASE_URL='postgresql+psycopg2://user:password@database:5432/castdb' \ +opensampl maintenance value-backfill +``` + +## Tune the value backfill + +```bash +opensampl maintenance value-backfill \ + --workers 4 \ + --batch-size 12h \ + --vacuum-every-batches 8 +``` + +The options are: + +- `--workers INTEGER`: number of concurrent database workers. If omitted, the + command uses `WORKERS`, then the local CPU count, then `4` as a fallback. +- `--batch-size DURATION`: time covered by each transaction. The default is + `1d`. Positive minute, hour, day, and week values are accepted, such as + `30m`, `12h`, `1d`, and `2w`. +- `--vacuum-every-batches INTEGER`: run `VACUUM` after this many completed + batches. The default is twice the resolved worker count. Use `0` to process + all batches first and vacuum once at the end. + +Smaller time windows reduce the amount of work lost if a transaction fails but +create more transactions. More workers may finish sooner, but increase database +CPU, I/O, connection usage, and write-ahead log activity. + +## Run in the packaged Compose deployment + +The packaged stack defines an opt-in `value-backfill` service. It uses the same +database and migration image as the rest of the deployment and is excluded from +normal `opensampl-server up` operations. + +After starting or upgrading the deployment, run it explicitly: + +```bash +opensampl-server run value-backfill +``` + +The service waits for a healthy database and successful migration completion. +It sets `ROUTE_TO_BACKEND=false` and supplies the container's direct +`DATABASE_URL`. + +To override the defaults, replace the service command while retaining its +environment and dependencies: + +```bash +opensampl-server run -- value-backfill \ + opensampl maintenance value-backfill \ + --workers 4 \ + --batch-size 12h \ + --vacuum-every-batches 8 +``` + +## Monitor and recover + +Each committed window is logged with its start time, end time, and updated row +count. Vacuum operations and the final total are also logged. A configuration, +schema, database, or worker error produces a nonzero exit status. + +If the value backfill fails, rerun it after correcting the error. You may retain +the same batch settings or lower the worker count to reduce database load. +Committed windows remain committed, and populated `value_float` rows are +skipped, so no manual checkpoint or cleanup is required. Running the command +after completion is safe; it exits when no values remain to be backfilled. diff --git a/mkdocs.yaml b/mkdocs.yaml index c439a1f..0b11d41 100644 --- a/mkdocs.yaml +++ b/mkdocs.yaml @@ -13,6 +13,7 @@ nav: - Expected Table Format: guides/expected_table_format.md - Create: guides/create_probe_type.md - Server: guides/opensampl-server.md + - Value Backfill: guides/value-backfill.md - Collect: guides/collection.md - Automatic Ingest: guides/automatic_ingest.md - NTP Extension: guides/ntp-extension.md diff --git a/opensampl/cli.py b/opensampl/cli.py index 3137c44..d53d1be 100644 --- a/opensampl/cli.py +++ b/opensampl/cli.py @@ -18,6 +18,7 @@ from opensampl.config.base import BaseConfig as CLIConfig from opensampl.db.orm import get_table_names +from opensampl.helpers.convert_value import value_backfill from opensampl.load_data import create_new_tables, write_to_table from opensampl.mixins.collect import CollectMixin from opensampl.mixins.random_data import RandomDataMixin @@ -176,6 +177,14 @@ def config_set(ctx: click.Context, name: str, value: str): conf.set_by_name(name=name, value=value) +@cli.group(cls=CaseInsensitiveGroup) +def maintenance(): + """Run direct-database maintenance operations.""" + + +maintenance.add_command(value_backfill) + + @cli.group(cls=CaseInsensitiveGroup) def load(): """Load data into database""" diff --git a/opensampl/config/server.py b/opensampl/config/server.py index 16d061b..3a8afbe 100644 --- a/opensampl/config/server.py +++ b/opensampl/config/server.py @@ -133,5 +133,5 @@ def get_db_url(self): password = self.docker_env_values.get("POSTGRES_PASSWORD") db = self.docker_env_values.get("POSTGRES_DB") if all(x is not None for x in [user, password, db]): - return f"postgresql://{user}:{password}@localhost:5415/{db}" + return f"postgresql+psycopg2://{user}:{password}@localhost:5415/{db}" raise ValueError("Database environment variables POSTGRES_USER, POSTGRES_PASSWORD, or POSTGRES_DB are not set.") diff --git a/opensampl/helpers/convert_value.py b/opensampl/helpers/convert_value.py new file mode 100644 index 0000000..a2e121f --- /dev/null +++ b/opensampl/helpers/convert_value.py @@ -0,0 +1,253 @@ +"""Backfill numeric probe values after the probe-data type migration.""" + +from __future__ import annotations + +import os +import re +from concurrent.futures import ThreadPoolExecutor +from datetime import datetime, timedelta +from typing import TYPE_CHECKING, Any + +import click +from loguru import logger +from sqlalchemy import create_engine, text + +if TYPE_CHECKING: + from collections.abc import Callable, Iterable + + from sqlalchemy import Engine + +DEFAULT_BATCH_SIZE = "1d" +_DURATION_PATTERN = re.compile(r"^(?P(?:\d+(?:\.\d*)?|\.\d+))(?P[mhdw])$", re.IGNORECASE) +_DURATION_SECONDS = { + "m": 60, + "h": 60 * 60, + "d": 24 * 60 * 60, + "w": 7 * 24 * 60 * 60, +} + +METRIC_TYPES_SQL = text(""" + SELECT uuid + FROM castdb.metric_type + WHERE value_type IN ('float', 'int') +""") + +SCHEMA_COLUMNS_SQL = text(""" + SELECT column_name + FROM information_schema.columns + WHERE table_schema = 'castdb' + AND table_name = 'probe_data' + AND column_name IN ('value_jsonb', 'value_float') +""") + +RANGE_SQL = text(""" + SELECT + min(pd.time) AS start_time, + max(pd.time) AS end_time + FROM castdb.probe_data pd + WHERE pd.metric_type_uuid = ANY(:metric_type_uuids) + AND pd.value_float IS NULL + AND pd.value_jsonb IS NOT NULL +""") + +UPDATE_SQL = text(""" + UPDATE castdb.probe_data + SET value_float = (value_jsonb #>> '{}')::double precision + WHERE metric_type_uuid = ANY(:metric_type_uuids) + AND value_float IS NULL + AND value_jsonb IS NOT NULL + AND time >= :start_time + AND time < :end_time +""") + + +def parse_duration(value: str) -> timedelta: + """Parse a positive duration expressed in minutes, hours, days, or weeks.""" + match = _DURATION_PATTERN.fullmatch(value.strip()) + if match is None: + raise click.BadParameter("must be a positive duration such as 30m, 12h, 1d, or 2w") + + seconds = float(match.group("value")) * _DURATION_SECONDS[match.group("unit").lower()] + if seconds <= 0: + raise click.BadParameter("must be greater than zero") + return timedelta(seconds=seconds) + + +def resolve_workers(workers: int | None) -> int: + """Resolve worker count from the option, environment, or local CPU count.""" + if workers is not None: + return workers + + env_workers = os.getenv("WORKERS") + if env_workers is not None: + try: + resolved = int(env_workers) + except ValueError as exc: + raise click.BadParameter("WORKERS must be a positive integer", param_hint="--workers") from exc + if resolved <= 0: + raise click.BadParameter("WORKERS must be a positive integer", param_hint="--workers") + return resolved + + return os.cpu_count() or 4 + + +def build_windows( + start_time: datetime, + end_time: datetime, + batch_size: timedelta, +) -> list[tuple[datetime, datetime]]: + """Build newest-first half-open time windows that include the final row.""" + end_exclusive = end_time + timedelta(microseconds=1) + windows: list[tuple[datetime, datetime]] = [] + current_start = start_time + while current_start < end_exclusive: + current_end = min(current_start + batch_size, end_exclusive) + windows.append((current_start, current_end)) + current_start = current_end + windows.reverse() + return windows + + +def _validate_schema(engine: Engine) -> None: + """Ensure the value-type migration has created both required columns.""" + with engine.connect() as conn: + columns = set(conn.execute(SCHEMA_COLUMNS_SQL).scalars().all()) + missing = {"value_jsonb", "value_float"} - columns + if missing: + missing_names = ", ".join(sorted(missing)) + raise RuntimeError( + f"The probe-data type migration has not been applied; missing castdb.probe_data column(s): {missing_names}" + ) + + +def _iter_waves( + windows: list[tuple[datetime, datetime]], + vacuum_every_batches: int, +) -> Iterable[list[tuple[datetime, datetime]]]: + """Split windows at vacuum boundaries, or return one wave when disabled.""" + wave_size = vacuum_every_batches or len(windows) + for index in range(0, len(windows), wave_size): + yield windows[index : index + wave_size] + + +def run_backfill( + database_url: str, + workers: int, + batch_size: timedelta, + vacuum_every_batches: int, + *, + engine_factory: Callable[..., Engine] | None = None, +) -> int: + """Backfill numeric values and return the number of rows updated.""" + factory = engine_factory or create_engine + engine = factory(database_url, pool_size=workers + 1) + logger.info("Using database: {}", engine.url) + logger.info("Workers: {}", workers) + + try: + _validate_schema(engine) + with engine.connect() as conn: + metric_type_uuids = conn.execute(METRIC_TYPES_SQL).scalars().all() + + if not metric_type_uuids: + logger.info("No float/int metric types found; nothing to map.") + return 0 + + with engine.connect() as conn: + time_range = conn.execute( + RANGE_SQL, + {"metric_type_uuids": metric_type_uuids}, + ).one() + + if time_range.start_time is None or time_range.end_time is None: + logger.info("No values require backfilling.") + return 0 + + windows = build_windows(time_range.start_time, time_range.end_time, batch_size) + logger.info("Backfilling values from {} through {}", time_range.start_time, time_range.end_time) + + def run_batch(window: tuple[datetime, datetime]) -> int: + batch_start, batch_end = window + with engine.begin() as conn: + conn.execute(text("SET LOCAL synchronous_commit = off")) + result = conn.execute( + UPDATE_SQL, + { + "metric_type_uuids": metric_type_uuids, + "start_time": batch_start, + "end_time": batch_end, + }, + ) + logger.info("{} -> {}: committed {:,} rows", batch_start, batch_end, result.rowcount) + return result.rowcount + + total_updated = 0 + vacuum_engine = engine.execution_options(isolation_level="AUTOCOMMIT") + with ThreadPoolExecutor(max_workers=workers) as pool: + for wave in _iter_waves(windows, vacuum_every_batches): + total_updated += sum(pool.map(run_batch, wave)) + logger.info("Vacuuming probe_data...") + with vacuum_engine.connect() as conn: + conn.execute(text("VACUUM castdb.probe_data")) + + logger.info("Value backfill complete; updated {:,} rows.", total_updated) + return total_updated + finally: + engine.dispose() + + +@click.command("value-backfill") +@click.option( + "--workers", + type=click.IntRange(min=1), + default=None, + help="Parallel database workers. Defaults to WORKERS or the local CPU count.", +) +@click.option( + "--batch-size", + type=str, + default=DEFAULT_BATCH_SIZE, + show_default=True, + metavar="DURATION", + help="Time span per batch, using m, h, d, or w units.", +) +@click.option( + "--vacuum-every-batches", + type=click.IntRange(min=0), + default=None, + metavar="INTEGER", + help="Vacuum after this many batches; defaults to workers * 2. Use 0 for only a final vacuum.", +) +@click.pass_obj +def value_backfill( + obj: dict[str, Any], + workers: int | None, + batch_size: str, + vacuum_every_batches: int | None, +) -> None: + """Map numeric JSON values into the optimized floating-point column.""" + config = obj["conf"] + if config.ROUTE_TO_BACKEND: + click.echo( + "Warning: value-backfill requires a direct database connection and cannot run when ROUTE_TO_BACKEND=true.", + err=True, + ) + raise click.UsageError("Disable ROUTE_TO_BACKEND before running this command.") + if not config.DATABASE_URL: + raise click.UsageError("DATABASE_URL must be configured to run value-backfill.") + + resolved_workers = resolve_workers(workers) + resolved_batch_size = parse_duration(batch_size) + resolved_vacuum_cadence = resolved_workers * 2 if vacuum_every_batches is None else vacuum_every_batches + + try: + run_backfill( + config.DATABASE_URL, + resolved_workers, + resolved_batch_size, + resolved_vacuum_cadence, + ) + except click.ClickException: + raise + except Exception as exc: + raise click.ClickException(f"Value backfill failed: {exc}") from exc diff --git a/opensampl/load_data.py b/opensampl/load_data.py index b4717bf..e8b9827 100644 --- a/opensampl/load_data.py +++ b/opensampl/load_data.py @@ -142,7 +142,7 @@ def load_time_data( # Ensure correct time dtypes df["time"] = pd.to_datetime(df["time"], format="mixed", utc=True, errors="raise") - if data_definition.metric.is_numeric(): + if metric_type.is_numeric(): df["value_float"] = pd.to_numeric(df["value"], errors="raise") df["value_jsonb"] = None else: diff --git a/opensampl/server/backend/main.py b/opensampl/server/backend/main.py index 87d3783..341b9c1 100644 --- a/opensampl/server/backend/main.py +++ b/opensampl/server/backend/main.py @@ -53,6 +53,9 @@ class ProbeMetadataPayload(BaseModel): DATABASE_URI = os.getenv("DATABASE_URL") +if DATABASE_URI and DATABASE_URI.startswith("postgresql://"): + DATABASE_URI = DATABASE_URI.replace("postgresql://", "postgresql+psycopg2://", 1) + engine = create_engine(DATABASE_URI) loglevel = os.getenv("BACKEND_LOG_LEVEL", "INFO") @@ -93,7 +96,7 @@ def get_keys(): Session = sessionmaker(bind=engine) # noqa: N806 with Session() as session: now = datetime.now(tz=UTC) - stmt = select(APIAccessKey.key).where(or_(APIAccessKey.expires_at is None, APIAccessKey.expires_at > now)) + stmt = select(APIAccessKey.key).where(or_(APIAccessKey.expires_at.is_(None), APIAccessKey.expires_at > now)) result = session.execute(stmt) keys = [row[0] for row in result.all()] logger.debug("api access keys loaded from db") diff --git a/opensampl/server/docker-bake.hcl b/opensampl/server/docker-bake.hcl new file mode 100644 index 0000000..8833f45 --- /dev/null +++ b/opensampl/server/docker-bake.hcl @@ -0,0 +1,55 @@ +// Docker Bake file for building and pushing the OpenSAMPL server images +// (backend + migrations) as multi-arch images tagged with both `latest` +// and the current opensampl package version. +// +// Usage: +// docker buildx bake -f docker-bake.hcl --push +// +// Override the registry or version if needed: +// REGISTRY=myregistry.example.com/opensampl OPENSAMPL_VERSION=1.2.1 \ +// docker buildx bake -f docker-bake.hcl --push + +variable "REGISTRY" { + default = "savannah.ornl.gov/opensampl" +} + +// Read the version straight out of the repo's pyproject.toml so the bake +// file never drifts from the package version. +variable "OPENSAMPL_VERSION" { + default = "1.2.3" +} + +variable "PLATFORMS" { + default = ["linux/amd64", "linux/arm64"] +} + +group "default" { + targets = ["backend", "migrations"] +} + +target "backend" { + context = "./backend" + dockerfile = "Dockerfile" + target = "prod" + platforms = PLATFORMS + args = { + OPENSAMPL_VERSION = OPENSAMPL_VERSION + } + tags = [ + "${REGISTRY}/backend:latest", + "${REGISTRY}/backend:${OPENSAMPL_VERSION}", + ] +} + +target "migrations" { + context = "./migrations" + dockerfile = "Dockerfile" + platforms = PLATFORMS + args = { + OPENSAMPL_VERSION = OPENSAMPL_VERSION + } + tags = [ + "${REGISTRY}/migrations:latest", + "${REGISTRY}/migrations:${OPENSAMPL_VERSION}", + ] +} diff --git a/opensampl/server/docker-compose.dev.yaml b/opensampl/server/docker-compose.dev.yaml index c1cc204..90cfa7b 100644 --- a/opensampl/server/docker-compose.dev.yaml +++ b/opensampl/server/docker-compose.dev.yaml @@ -44,11 +44,27 @@ services: context: migrations restart: "no" environment: - - DB_URI=postgresql://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB} + - DB_URI=postgresql+psycopg2://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB} depends_on: db: condition: service_healthy + value-backfill: + image: savannah.ornl.gov/opensampl/migrations:latest + build: + context: migrations + profiles: ["maintenance"] + restart: "no" + command: ["opensampl", "maintenance", "value-backfill"] + environment: + - DATABASE_URL=postgresql+psycopg2://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB} + - ROUTE_TO_BACKEND=false + depends_on: + db: + condition: service_healthy + migrations: + condition: service_completed_successfully + backend: image: savannah.ornl.gov/opensampl/backend:latest @@ -59,7 +75,7 @@ services: - "8015:8000" restart: unless-stopped environment: - - DATABASE_URL=postgresql://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB} + - DATABASE_URL=postgresql+psycopg2://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB} - ROUTE_TO_BACKEND=false - BACKEND_LOG_LEVEL=${BACKEND_LOG_LEVEL} - USE_API_KEY=${USE_API_KEY} @@ -72,4 +88,4 @@ services: volumes: castdb: - grafana-data: \ No newline at end of file + grafana-data: diff --git a/opensampl/server/docker-compose.yaml b/opensampl/server/docker-compose.yaml index 48fa367..5a7184b 100644 --- a/opensampl/server/docker-compose.yaml +++ b/opensampl/server/docker-compose.yaml @@ -40,13 +40,27 @@ services: context: ./migrations restart: "no" environment: - - DB_URI=postgresql://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB} + - DB_URI=postgresql+psycopg2://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB} - CHUNK_INTERVAL=${CHUNK_INTERVAL} - RETENTION_POLICY=${RETENTION_POLICY} depends_on: db: condition: service_healthy + value-backfill: + image: savannah.ornl.gov/opensampl/migrations:latest + profiles: ["maintenance"] + restart: "no" + command: ["opensampl", "maintenance", "value-backfill"] + environment: + - DATABASE_URL=postgresql+psycopg2://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB} + - ROUTE_TO_BACKEND=false + depends_on: + db: + condition: service_healthy + migrations: + condition: service_completed_successfully + backend: image: savannah.ornl.gov/opensampl/backend:latest build: @@ -56,7 +70,7 @@ services: - "8015:8000" restart: always environment: - - DATABASE_URL=postgresql://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB} + - DATABASE_URL=postgresql+psycopg2://${POSTGRES_USER}:${POSTGRES_PASSWORD}@db:5432/${POSTGRES_DB} - BACKEND_LOG_LEVEL=${BACKEND_LOG_LEVEL:-INFO} - USE_API_KEY=${USE_API_KEY:-false} - API_KEY=${API_KEY:-} @@ -66,4 +80,4 @@ services: volumes: castdb: - grafana-data: \ No newline at end of file + grafana-data: diff --git a/opensampl/server/migrations/Dockerfile b/opensampl/server/migrations/Dockerfile index d3c7ae1..28c9378 100644 --- a/opensampl/server/migrations/Dockerfile +++ b/opensampl/server/migrations/Dockerfile @@ -3,7 +3,7 @@ FROM python:3.12 ARG OPENSAMPL_VERSION=1.1.5 WORKDIR / -RUN pip install --no-cache-dir "opensampl[migrations]==${OPENSAMPL_VERSION}" alembic +RUN pip install --no-cache-dir "opensampl[migrations]==${OPENSAMPL_VERSION}" RUN useradd -m alembic USER alembic diff --git a/opensampl/server/migrations/_migrations/env.py b/opensampl/server/migrations/_migrations/env.py index a44ba86..b8d9ab1 100644 --- a/opensampl/server/migrations/_migrations/env.py +++ b/opensampl/server/migrations/_migrations/env.py @@ -6,7 +6,11 @@ from alembic import context -sqlalchemy_url = os.environ.get('DB_URI') +db_uri = os.environ.get('DB_URI') +if db_uri and db_uri.startswith('postgresql://'): + db_uri = db_uri.replace('postgresql://', 'postgresql+psycopg2://', 1) + +sqlalchemy_url = db_uri # this is the Alembic Config object, which provides # access to the values within the .ini file in use. config = context.config diff --git a/opensampl/server/migrations/_migrations/versions/2026_09_11_1445_update_probe_data_types.py b/opensampl/server/migrations/_migrations/versions/2026_09_11_1445_update_probe_data_types.py index 7c2688d..74f50bb 100644 --- a/opensampl/server/migrations/_migrations/versions/2026_09_11_1445_update_probe_data_types.py +++ b/opensampl/server/migrations/_migrations/versions/2026_09_11_1445_update_probe_data_types.py @@ -30,18 +30,6 @@ def upgrade() -> None: sa.Column("value_float", sa.Float(), nullable=True), schema=SCHEMA) - # Backfill value_float from the numeric jsonb values - op.execute( - """ - UPDATE castdb.probe_data pd - SET value_float = (pd.value_jsonb #>> '{}')::double precision - FROM castdb.metric_type mt - WHERE pd.metric_type_uuid = mt.uuid - AND mt.value_type IN ('float' - , 'int') - """ - ) - def downgrade() -> None: op.execute( diff --git a/pyproject.toml b/pyproject.toml index 55369d2..14f2dc4 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "opensampl" -version = "1.2.2" +version = "1.2.3" description = "Python tools for adding clock data to a timescale db." license = {file = "LICENSE"} authors = [ @@ -43,7 +43,7 @@ dependencies = [ "pydantic>=2.10.3,<3", "pydantic-settings>=2.9.0", "sqlalchemy>=2.0.39,<3", - "geoalchemy2==0.16.0", + "geoalchemy2>=0.18.0,<1", "click>=8.0.0,<9", "pandas>=2.2.1,<3", "tqdm>=4.66.2,<5", diff --git a/tests/test_cli.py b/tests/test_cli.py index f5d5e46..e7b401a 100644 --- a/tests/test_cli.py +++ b/tests/test_cli.py @@ -177,6 +177,22 @@ def test_cli_init_command(self, runner): assert result.exit_code == 0 + def test_cli_maintenance_command(self, runner): + """Test the maintenance command group.""" + result = runner.invoke(cli, ["maintenance", "--help"]) + + assert result.exit_code == 0 + assert "value-backfill" in result.output + + def test_cli_value_backfill_command(self, runner): + """Test the value-backfill command help.""" + result = runner.invoke(cli, ["maintenance", "value-backfill", "--help"]) + + assert result.exit_code == 0 + assert "--workers" in result.output + assert "--batch-size" in result.output + assert "--vacuum-every-batches" in result.output + def test_cli_case_insensitive_commands(self, runner): """Test case-insensitive subcommand handling for 'load'.""" # Only subcommands of 'load' are case-insensitive, not the top-level @@ -185,4 +201,5 @@ def test_cli_case_insensitive_commands(self, runner): result3 = runner.invoke(cli, ['load', 'Table', '--help']) # All should work the same - assert result1.exit_code == result2.exit_code == result3.exit_code == 0 \ No newline at end of file + assert result1.exit_code == result2.exit_code == result3.exit_code == 0 + diff --git a/tests/test_config.py b/tests/test_config.py index 74bc6b2..c7f44d6 100644 --- a/tests/test_config.py +++ b/tests/test_config.py @@ -194,7 +194,7 @@ def test_get_db_url(self): } db_url = config.get_db_url() - assert db_url == "postgresql://testuser:testpass@localhost:5415/testdb" + assert db_url == "postgresql+psycopg2://testuser:testpass@localhost:5415/testdb" # Test missing environment variables config.docker_env_values = {} diff --git a/tests/test_convert_value.py b/tests/test_convert_value.py new file mode 100644 index 0000000..c0b5678 --- /dev/null +++ b/tests/test_convert_value.py @@ -0,0 +1,273 @@ +"""Tests for the numeric value maintenance backfill.""" + +from datetime import datetime, timedelta +from unittest.mock import MagicMock, Mock, patch + +import click +import pytest +from click.testing import CliRunner + +from opensampl.cli import cli +from opensampl.helpers.convert_value import build_windows, parse_duration, resolve_workers, run_backfill + + +@pytest.mark.parametrize( + ("value", "expected"), + [ + ("30m", timedelta(minutes=30)), + ("12h", timedelta(hours=12)), + ("1d", timedelta(days=1)), + ("2w", timedelta(weeks=2)), + ("1.5h", timedelta(minutes=90)), + ], +) +def test_parse_duration(value, expected): + """Supported duration units are converted to timedeltas.""" + assert parse_duration(value) == expected + + +@pytest.mark.parametrize("value", ["", "0h", "-1d", "1", "tomorrow"]) +def test_parse_duration_rejects_invalid_values(value): + """Invalid and nonpositive durations are rejected.""" + with pytest.raises(click.BadParameter): + parse_duration(value) + + +def test_resolve_workers_prefers_argument(monkeypatch): + """An explicit worker count takes precedence over the environment.""" + monkeypatch.setenv("WORKERS", "9") + + assert resolve_workers(3) == 3 + + +def test_resolve_workers_uses_environment(monkeypatch): + """WORKERS supplies the default when the option is omitted.""" + monkeypatch.setenv("WORKERS", "7") + + assert resolve_workers(None) == 7 + + +def test_build_windows_is_newest_first_and_includes_end(): + """Batch windows are returned newest-first with an exclusive final bound.""" + start = datetime(2026, 1, 1) + end = datetime(2026, 1, 3) + + windows = build_windows(start, end, timedelta(days=1)) + + assert windows == [ + (datetime(2026, 1, 3), datetime(2026, 1, 3, 0, 0, 0, 1)), + (datetime(2026, 1, 2), datetime(2026, 1, 3)), + (datetime(2026, 1, 1), datetime(2026, 1, 2)), + ] + + +def test_route_to_backend_warns_and_aborts_before_connecting(tmp_path): + """The maintenance operation cannot bypass configured backend routing.""" + env_file = tmp_path / ".env" + env_file.write_text( + "ROUTE_TO_BACKEND=true\n" + "BACKEND_URL=http://backend:8000\n" + "DATABASE_URL=postgresql://db/test\n" + ) + + with patch("opensampl.helpers.convert_value.run_backfill") as mock_backfill: + result = CliRunner().invoke( + cli, + ["--env-file", str(env_file), "maintenance", "value-backfill"], + ) + + assert result.exit_code != 0 + assert "Warning:" in result.output + assert "Disable ROUTE_TO_BACKEND" in result.output + mock_backfill.assert_not_called() + + +def test_map_values_requires_database_url(tmp_path): + """A direct database URL is required.""" + env_file = tmp_path / ".env" + env_file.write_text("ROUTE_TO_BACKEND=false\n") + + result = CliRunner().invoke( + cli, + ["--env-file", str(env_file), "maintenance", "value-backfill"], + ) + + assert result.exit_code != 0 + assert "DATABASE_URL must be configured" in result.output + + +def test_map_values_forwards_resolved_options(tmp_path): + """CLI options are resolved and forwarded to the backfill implementation.""" + env_file = tmp_path / ".env" + env_file.write_text("ROUTE_TO_BACKEND=false\nDATABASE_URL=postgresql://db/test\n") + + with patch("opensampl.helpers.convert_value.run_backfill", return_value=0) as mock_backfill: + result = CliRunner().invoke( + cli, + [ + "--env-file", + str(env_file), + "maintenance", + "value-backfill", + "--workers", + "3", + "--batch-size", + "12h", + "--vacuum-every-batches", + "0", + ], + ) + + assert result.exit_code == 0 + mock_backfill.assert_called_once_with( + "postgresql://db/test", + 3, + timedelta(hours=12), + 0, + ) + + +def test_map_values_derives_vacuum_cadence(tmp_path): + """Vacuum cadence defaults to twice the resolved worker count.""" + env_file = tmp_path / ".env" + env_file.write_text("ROUTE_TO_BACKEND=false\nDATABASE_URL=postgresql://db/test\n") + + with patch("opensampl.helpers.convert_value.run_backfill", return_value=0) as mock_backfill: + result = CliRunner().invoke( + cli, + [ + "--env-file", + str(env_file), + "maintenance", + "value-backfill", + "--workers", + "5", + ], + ) + + assert result.exit_code == 0 + assert mock_backfill.call_args.args[3] == 10 + + +def test_run_backfill_stops_when_no_metric_types_exist(): + """An empty metric selection exits without creating batches or vacuuming.""" + engine = Mock() + engine.url = "postgresql://db/test" + + schema_connection = Mock() + schema_connection.execute.return_value.scalars.return_value.all.return_value = ["value_jsonb", "value_float"] + schema_context = MagicMock() + schema_context.__enter__.return_value = schema_connection + + metric_connection = Mock() + metric_connection.execute.return_value.scalars.return_value.all.return_value = [] + metric_context = MagicMock() + metric_context.__enter__.return_value = metric_connection + + engine.connect.side_effect = [schema_context, metric_context] + engine_factory = Mock(return_value=engine) + + updated = run_backfill( + "postgresql://db/test", + workers=2, + batch_size=timedelta(days=1), + vacuum_every_batches=4, + engine_factory=engine_factory, + ) + + assert updated == 0 + engine.execution_options.assert_not_called() + engine.dispose.assert_called_once_with() + + +def test_run_backfill_updates_windows_and_vacuums_once(): + """A zero cadence processes all windows before one final vacuum.""" + engine = Mock() + engine.url = "postgresql://db/test" + + schema_connection = Mock() + schema_connection.execute.return_value.scalars.return_value.all.return_value = ["value_jsonb", "value_float"] + schema_context = MagicMock() + schema_context.__enter__.return_value = schema_connection + + metric_connection = Mock() + metric_connection.execute.return_value.scalars.return_value.all.return_value = ["metric-uuid"] + metric_context = MagicMock() + metric_context.__enter__.return_value = metric_connection + + time_range = Mock(start_time=datetime(2026, 1, 1), end_time=datetime(2026, 1, 2)) + range_connection = Mock() + range_connection.execute.return_value.one.return_value = time_range + range_context = MagicMock() + range_context.__enter__.return_value = range_connection + engine.connect.side_effect = [schema_context, metric_context, range_context] + + transaction_contexts = [] + for row_count in (2, 3): + transaction_connection = Mock() + transaction_connection.execute.side_effect = [Mock(), Mock(rowcount=row_count)] + transaction_context = MagicMock() + transaction_context.__enter__.return_value = transaction_connection + transaction_contexts.append(transaction_context) + engine.begin.side_effect = transaction_contexts + + vacuum_connection = Mock() + vacuum_context = MagicMock() + vacuum_context.__enter__.return_value = vacuum_connection + vacuum_engine = Mock() + vacuum_engine.connect.return_value = vacuum_context + engine.execution_options.return_value = vacuum_engine + + updated = run_backfill( + "postgresql://db/test", + workers=1, + batch_size=timedelta(days=1), + vacuum_every_batches=0, + engine_factory=Mock(return_value=engine), + ) + + assert updated == 5 + assert engine.begin.call_count == 2 + vacuum_connection.execute.assert_called_once() + engine.dispose.assert_called_once_with() + + +def test_run_backfill_propagates_batch_failure_and_disposes_engine(): + """A failed batch aborts the run while still disposing the engine.""" + engine = Mock() + engine.url = "postgresql://db/test" + + schema_connection = Mock() + schema_connection.execute.return_value.scalars.return_value.all.return_value = ["value_jsonb", "value_float"] + schema_context = MagicMock() + schema_context.__enter__.return_value = schema_connection + + metric_connection = Mock() + metric_connection.execute.return_value.scalars.return_value.all.return_value = ["metric-uuid"] + metric_context = MagicMock() + metric_context.__enter__.return_value = metric_connection + + time_range = Mock(start_time=datetime(2026, 1, 1), end_time=datetime(2026, 1, 1, 1)) + range_connection = Mock() + range_connection.execute.return_value.one.return_value = time_range + range_context = MagicMock() + range_context.__enter__.return_value = range_connection + engine.connect.side_effect = [schema_context, metric_context, range_context] + + transaction_connection = Mock() + transaction_connection.execute.side_effect = [Mock(), RuntimeError("batch failed")] + transaction_context = MagicMock() + transaction_context.__enter__.return_value = transaction_connection + engine.begin.return_value = transaction_context + engine.execution_options.return_value = Mock() + + with pytest.raises(RuntimeError, match="batch failed"): + run_backfill( + "postgresql://db/test", + workers=1, + batch_size=timedelta(days=1), + vacuum_every_batches=2, + engine_factory=Mock(return_value=engine), + ) + + engine.dispose.assert_called_once_with() diff --git a/uv.lock b/uv.lock index 5fd41d7..4f1a872 100644 --- a/uv.lock +++ b/uv.lock @@ -443,15 +443,15 @@ wheels = [ [[package]] name = "geoalchemy2" -version = "0.16.0" +version = "0.20.0" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "packaging" }, { name = "sqlalchemy" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/b5/d6/b01fa413cb8d7b3407f7f832f9310c623597ba38d21e52082296a36c04fa/geoalchemy2-0.16.0.tar.gz", hash = "sha256:df64bb72af70daafaac3f359492c96501c37ab85ed20f9510c99cc6d02881100", size = 227842, upload-time = "2024-11-13T08:24:37.98Z" } +sdist = { url = "https://files.pythonhosted.org/packages/63/74/6cb1ef591bf47d28f41aa770f2f3a91c0a570aee0a4083bed7f8c533d8df/geoalchemy2-0.20.0.tar.gz", hash = "sha256:450f427f4bc3cf2d5ddee0af3763aed0f3eea2384e7c9a99798d8f1508279322", size = 280805, upload-time = "2026-05-12T14:50:26.132Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/12/43/b1a687122767d544d493776934687ab168b2224bd4a6c1e8b22640ce6227/GeoAlchemy2-0.16.0-py3-none-any.whl", hash = "sha256:b0f27d5500ee757af4654c6262e0f834b7a843504d193653ec747ef1128d2ab5", size = 74663, upload-time = "2024-11-13T08:24:36.234Z" }, + { url = "https://files.pythonhosted.org/packages/4e/08/b66ad4239f592e05202e25925c08cdd04cc14c3994000ec70ec61fea202c/geoalchemy2-0.20.0-py3-none-any.whl", hash = "sha256:1489a1d106519542a79c97cd0b4c537d80462c353610ebc2429cf2c43daac717", size = 96467, upload-time = "2026-05-12T14:50:24.998Z" }, ] [[package]] @@ -1090,7 +1090,7 @@ wheels = [ [[package]] name = "opensampl" -version = "1.2.0" +version = "1.2.3" source = { editable = "." } dependencies = [ { name = "allantools" }, @@ -1153,7 +1153,7 @@ requires-dist = [ { name = "astor" }, { name = "click", specifier = ">=8.0.0,<9" }, { name = "fastapi", marker = "extra == 'backend'" }, - { name = "geoalchemy2", specifier = "==0.16.0" }, + { name = "geoalchemy2", specifier = ">=0.18.0,<1" }, { name = "jinja2", specifier = ">=3.1.6" }, { name = "libcst" }, { name = "loguru", specifier = ">=0.7.0,<0.8" },