Source code for qcodes.dataset._raw_data_storage

"""
Module for managing per-dataset raw data SQLite files.

When ``dataset.raw_data_backend`` is set to ``"sqlite_per_dataset_db"``,
measurement data (results tables) are written to individual SQLite files
- one per dataset - instead of the main QCoDeS database file.  All metadata
(runs, experiments, parameters) remains in the main database.

The per-dataset files are stored in the folder given by
``dataset.raw_data_backend_config.sqlite_per_dataset_db.raw_data_path`` and are
named ``<guid>.db``.
"""

from __future__ import annotations

import logging
from contextlib import closing
from dataclasses import dataclass, field
from datetime import UTC, datetime, timedelta
from pathlib import Path
from typing import TYPE_CHECKING, NamedTuple

from tqdm.auto import tqdm

import qcodes
from qcodes.dataset.export_config import _expand_export_path
from qcodes.dataset.sqlite.connection import AtomicConnection, atomic
from qcodes.dataset.sqlite.database import (
    _connect_to_sqlite_file,
    connect,
)
from qcodes.dataset.sqlite.queries import (
    _create_run_table,
    _remove_dataset_from_db,
    get_datasets_with_raw_data_path,
)
from qcodes.dataset.sqlite.query_helpers import insert_column, is_column_in_table

if TYPE_CHECKING:
    from collections.abc import Mapping, Sequence
    from typing import Any

    from qcodes.parameters import ParamSpecBase

log = logging.getLogger(__name__)

_DATASET_CONFIG_SECTION = "dataset"
_RESULTS_BACKEND_KEY = "raw_data_backend"
_BACKEND_CONFIG_KEY = "raw_data_backend_config"
_RAW_DATA_PATH_KEY = "raw_data_path"

#: ``raw_data_backend`` value keeping results in the main QCoDeS database.
MAIN_DB_BACKEND = "sqlite_main_db"
#: ``raw_data_backend`` value selecting the per-dataset SQLite-file backend.
PER_DATASET_DB_BACKEND = "sqlite_per_dataset_db"


def get_configured_results_backend_name() -> str:
    """Return the results backend name from ``dataset.raw_data_backend``."""
    return qcodes.config[_DATASET_CONFIG_SECTION].get(
        _RESULTS_BACKEND_KEY, MAIN_DB_BACKEND
    )


def get_results_backend_config(backend_name: str) -> Mapping[str, Any]:
    """Return the backend-specific config for *backend_name*.

    Reads ``dataset.raw_data_backend_config.<backend_name>``; returns an empty
    mapping when the backend has no configured settings.
    """
    all_config = qcodes.config[_DATASET_CONFIG_SECTION].get(_BACKEND_CONFIG_KEY, {})
    return all_config.get(backend_name, {}) or {}


def is_raw_data_storage_enabled() -> bool:
    """Return True if the per-dataset SQLite results backend is selected."""
    return get_configured_results_backend_name() == PER_DATASET_DB_BACKEND


def get_raw_data_folder(db_path: str | None = None) -> Path:
    """Return the resolved folder path for raw data SQLite files.

    The path template comes from the ``sqlite_per_dataset_db`` backend config
    and is expanded the same way as the export path. ``{db_location}`` is
    resolved relative to *db_path* (the owning dataset's database) when given,
    otherwise relative to the global ``core.db_location`` config.
    """
    config = get_results_backend_config(PER_DATASET_DB_BACKEND)
    raw_path_template: str = config.get(_RAW_DATA_PATH_KEY, "{db_location}")
    return Path(_expand_export_path(raw_path_template, db_path)).expanduser().absolute()


def get_raw_data_db_path(
    guid: str, folder: Path | None = None, db_path: str | None = None
) -> Path:
    """Return the full path for a dataset's raw data SQLite file.

    Args:
        guid: The GUID of the dataset.
        folder: Override folder.  If *None*, uses :func:`get_raw_data_folder`.
        db_path: The owning dataset's database file, used to resolve
            ``{db_location}`` when *folder* is not given.

    """
    if folder is None:
        folder = get_raw_data_folder(db_path)
    return folder / f"{guid}.db"


def connect_to_raw_data_db(
    path: str | Path,
    *,
    read_only: bool = False,
) -> AtomicConnection:
    """Open (or create) a lightweight SQLite connection for raw data.

    Unlike the main QCoDeS :func:`~qcodes.dataset.sqlite.database.connect`,
    this does **not** create the full metadata schema (experiments, runs, ...).
    It reuses the shared
    :func:`~qcodes.dataset.sqlite.database._connect_to_sqlite_file` helper so
    that the numpy/sqlite type adapters and connection settings are identical
    to the main database connection.

    Args:
        path: Path to the raw-data SQLite file.
        read_only: Open the database in read-only mode.

    Returns:
        An :class:`AtomicConnection` to the raw-data database.

    """
    return _connect_to_sqlite_file(path, read_only=read_only)


def create_raw_data_db(
    path: str | Path,
    table_name: str,
    paramspecs: Sequence[ParamSpecBase],
) -> AtomicConnection:
    """Create a per-dataset raw-data SQLite file with a results table.

    The file is created if it does not exist.  The parent directory is
    created if needed.

    Args:
        path: Full path for the new SQLite file.
        table_name: Name of the results table to create (matches the name
            in the main database).
        paramspecs: Parameter specifications describing the columns.

    Returns:
        An :class:`AtomicConnection` to the newly created database.

    """
    path = Path(path)
    path.parent.mkdir(parents=True, exist_ok=True)

    conn = connect_to_raw_data_db(path)

    # Create the empty results table, then add the parameter columns with
    # insert_column - exactly as the main database does - so column names are
    # quoted. ``_create_run_table`` with paramspecs would emit them verbatim via
    # ``ParamSpecBase.sql_repr()``, which breaks for names that are SQL keywords
    # (e.g. a parameter called ``from``).
    _create_run_table(conn, table_name)
    with atomic(conn) as aconn:
        for spec in paramspecs:
            insert_column(aconn, table_name, spec.name, spec.type)

    log.info(
        "Created raw data database at %s with table %s",
        path,
        table_name,
    )
    return conn


[docs] def update_raw_data_paths( db_path: str | Path, new_raw_data_folder: str | Path, ) -> list[tuple[int, str, str]]: """Update raw data file paths in the main database after files have moved. Use this when per-dataset raw data files have been relocated to a new folder but the main database still references the old paths. The function scans all runs that have a ``raw_data_db_path`` internal ``runs`` column set, verifies that a file with the expected GUID-based name exists in *new_raw_data_folder*, and updates the stored path in the database. Args: db_path: Path to the main QCoDeS database file. new_raw_data_folder: The new folder where the per-dataset SQLite files now reside. Returns: A list of ``(run_id, old_path, new_path)`` tuples for every run whose path was updated. Raises: FileNotFoundError: If the main database file does not exist. FileNotFoundError: If *new_raw_data_folder* does not exist. """ db_path = Path(db_path) # Resolve to an absolute path so the stored path does not depend on the # process working directory when datasets are loaded later (newly created # split datasets also store absolute paths). new_raw_data_folder = Path(new_raw_data_folder).resolve() if not db_path.is_file(): raise FileNotFoundError(f"Database file not found: {db_path}") if not new_raw_data_folder.is_dir(): raise FileNotFoundError(f"New raw data folder not found: {new_raw_data_folder}") updated: list[tuple[int, str, str]] = [] with closing(connect(str(db_path))) as conn: if not is_column_in_table(conn, "runs", "raw_data_db_path"): log.info( "No raw_data_db_path column found in %s; nothing to update.", db_path ) return [] cursor = conn.execute( "SELECT run_id, raw_data_db_path FROM runs " "WHERE raw_data_db_path IS NOT NULL" ) rows = cursor.fetchall() for run_id, old_path_str in tqdm(rows, desc="Updating raw data paths"): old_path = Path(old_path_str) # The per-dataset file name is always <guid>.db — preserved on move new_path = new_raw_data_folder / old_path.name if not new_path.is_file(): log.warning( "Run %d: expected raw data file %s not found in new folder;" " skipping.", run_id, new_path, ) continue if str(new_path) == old_path_str: continue # already correct new_path_str = str(new_path) with atomic(conn) as aconn: aconn.execute( "UPDATE runs SET raw_data_db_path = ? WHERE run_id = ?", (new_path_str, run_id), ) updated.append((run_id, old_path_str, new_path_str)) log.debug( "Run %d: updated raw_data_db_path from %s to %s", run_id, old_path_str, new_path_str, ) log.info("Updated %d raw data paths in %s", len(updated), db_path) return updated
# --------------------------------------------------------------------------- # Dataset management helpers # --------------------------------------------------------------------------- class DatasetInfo(NamedTuple): """Summary information about a dataset in the main database.""" run_id: int guid: str experiment_name: str sample_name: str run_timestamp: float | None completed_timestamp: float | None result_table_name: str raw_data_db_path: str | None raw_data_size_bytes: int | None @dataclass class PurgeResult: """Result of a purge_orphaned_datasets operation.""" total_datasets_with_raw_data: int orphaned_datasets: list[DatasetInfo] removed_datasets: list[DatasetInfo] dry_run: bool errors: list[tuple[int, Exception]] = field(default_factory=list) @dataclass class CleanupResult: """Result of a cleanup_datasets operation.""" total_datasets_scanned: int matching_datasets: list[DatasetInfo] removed_datasets: list[DatasetInfo] total_size_freed_bytes: int dry_run: bool errors: list[tuple[int, Exception]] = field(default_factory=list) def _build_dataset_info_list( conn: AtomicConnection, ) -> list[DatasetInfo]: """Query datasets with raw data paths and enrich with file size info.""" rows = get_datasets_with_raw_data_path(conn) datasets: list[DatasetInfo] = [] for row in tqdm(rows, desc="Scanning datasets"): raw_path = row.raw_data_db_path raw_size: int | None = None if raw_path and Path(raw_path).is_file(): raw_size = Path(raw_path).stat().st_size datasets.append( DatasetInfo( run_id=row.run_id, guid=row.guid, experiment_name=row.experiment_name, sample_name=row.sample_name, run_timestamp=row.run_timestamp, completed_timestamp=row.completed_timestamp, result_table_name=row.result_table_name, raw_data_db_path=raw_path, raw_data_size_bytes=raw_size, ) ) return datasets
[docs] def purge_orphaned_datasets( db_path: str | Path, *, dry_run: bool = True, ) -> PurgeResult: """Find and optionally remove dataset records whose raw data files are missing. When using split raw data storage, users may archive and delete individual per-dataset SQLite files. This function identifies datasets in the main database that reference raw data files which no longer exist on disk, and optionally removes those dataset records from the main database. Args: db_path: Path to the main QCoDeS database file. dry_run: If *True* (default), only report which datasets would be removed without making any changes. Set to *False* to actually delete the orphaned dataset records. Returns: A :class:`PurgeResult` with the list of orphaned datasets and, if *dry_run* is False, the list of datasets that were removed. Raises: FileNotFoundError: If the main database file does not exist. """ db_path = Path(db_path) if not db_path.is_file(): raise FileNotFoundError(f"Database file not found: {db_path}") with closing(connect(str(db_path))) as conn: all_datasets = _build_dataset_info_list(conn) # Orphaned = raw_data_size_bytes is None (file not found on disk) orphaned = [ds for ds in all_datasets if ds.raw_data_size_bytes is None] msg = ( f"Found {len(all_datasets)} datasets with raw data references in {db_path}, " f"{len(orphaned)} orphaned (file missing)." ) log.info(msg) print(msg) removed: list[DatasetInfo] = [] errors: list[tuple[int, Exception]] = [] if not dry_run and orphaned: for ds_info in tqdm(orphaned, desc="Removing orphaned datasets"): try: _remove_dataset_from_db(conn, ds_info.run_id) removed.append(ds_info) log.debug( "Removed orphaned dataset run_id=%d (guid=%s) from %s.", ds_info.run_id, ds_info.guid, db_path, ) except Exception as exc: # noqa: BLE001 collect any failure and continue purging the remaining datasets log.error("Failed to remove run_id=%d: %s", ds_info.run_id, exc) errors.append((ds_info.run_id, exc)) result = PurgeResult( total_datasets_with_raw_data=len(all_datasets), orphaned_datasets=orphaned, removed_datasets=removed, dry_run=dry_run, errors=errors, ) if dry_run: msg = f"Dry run: {len(orphaned)} orphaned datasets would be removed from {db_path}." else: msg = f"Removed {len(removed)} orphaned datasets from {db_path}." log.info(msg) print(msg) return result
[docs] def cleanup_datasets( db_path: str | Path, *, older_than_days: int | None = None, sample_name: str | None = None, larger_than_mb: float | None = None, dry_run: bool = True, ) -> CleanupResult: """Remove datasets and their raw data files matching given criteria. This function helps manage disk space by removing datasets that match one or more of the specified criteria. It removes both the raw data SQLite file on disk and the corresponding records in the main database. Criteria are combined with AND logic: a dataset must match **all** specified criteria to be selected for removal. Specify at least one criterion. In-progress runs (started but not yet completed) are never removed, since their raw-data file may still be actively written by another process. Args: db_path: Path to the main QCoDeS database file. older_than_days: Remove datasets whose *completed_timestamp* (or *run_timestamp* if not completed) is older than this many days ago. sample_name: Remove datasets belonging to experiments with this exact sample name. larger_than_mb: Remove datasets whose raw data file is larger than this many megabytes. dry_run: If *True* (default), only report which datasets would be removed without making any changes. Set to *False* to actually delete datasets and their raw data files. Returns: A :class:`CleanupResult` with details of the operation. Raises: FileNotFoundError: If the main database file does not exist. ValueError: If no criteria are specified, or if ``older_than_days`` or ``larger_than_mb`` is negative. """ db_path = Path(db_path) if not db_path.is_file(): raise FileNotFoundError(f"Database file not found: {db_path}") if older_than_days is None and sample_name is None and larger_than_mb is None: raise ValueError("At least one cleanup criterion must be specified.") # Reject negative thresholds: a negative age puts the cutoff in the future # and a negative size is below every file, so either would match (and, when # dry_run=False, delete) essentially all split datasets. if older_than_days is not None and older_than_days < 0: raise ValueError( f"older_than_days must be non-negative, got {older_than_days}." ) if larger_than_mb is not None and larger_than_mb < 0: raise ValueError(f"larger_than_mb must be non-negative, got {larger_than_mb}.") with closing(connect(str(db_path))) as conn: all_datasets = _build_dataset_info_list(conn) # Apply filters (AND logic) matching: list[DatasetInfo] = [] cutoff_ts: float | None = None if older_than_days is not None: cutoff_dt = datetime.now(tz=UTC) - timedelta(days=older_than_days) cutoff_ts = cutoff_dt.timestamp() size_threshold_bytes: int | None = None if larger_than_mb is not None: size_threshold_bytes = int(larger_than_mb * 1024 * 1024) for ds in all_datasets: # Never delete an in-progress run (started but not completed): # another process may still be writing its raw-data file, and # unlinking it would cause data loss. Such runs are skipped # regardless of the other criteria. if ds.completed_timestamp is None: continue # Age filter if cutoff_ts is not None: ts = ds.completed_timestamp or ds.run_timestamp if ts is None or ts >= cutoff_ts: continue # Sample name filter if sample_name is not None and ds.sample_name != sample_name: continue # Size filter if size_threshold_bytes is not None and ( ds.raw_data_size_bytes is None or ds.raw_data_size_bytes <= size_threshold_bytes ): continue matching.append(ds) msg = ( f"Found {len(all_datasets)} datasets with raw data in {db_path}, " f"{len(matching)} match cleanup criteria." ) log.info(msg) print(msg) removed: list[DatasetInfo] = [] errors: list[tuple[int, Exception]] = [] total_freed: int = 0 if not dry_run and matching: for ds_info in tqdm(matching, desc="Removing datasets"): try: raw_path = ( Path(ds_info.raw_data_db_path) if ds_info.raw_data_db_path else None ) # Remove the DB record first: if this fails, the raw data # file is still on disk (recoverable) rather than deleted # while its run keeps pointing at missing data. _remove_dataset_from_db(conn, ds_info.run_id) # Only after the record is gone, delete the raw data file. if raw_path is not None and raw_path.is_file(): file_size = raw_path.stat().st_size raw_path.unlink() total_freed += file_size log.debug("Deleted raw data file: %s", raw_path) removed.append(ds_info) log.debug( "Removed dataset run_id=%d (guid=%s) from %s.", ds_info.run_id, ds_info.guid, db_path, ) except Exception as exc: # noqa: BLE001 collect any failure and continue cleaning up the remaining datasets log.error("Failed to remove run_id=%d: %s", ds_info.run_id, exc) errors.append((ds_info.run_id, exc)) result = CleanupResult( total_datasets_scanned=len(all_datasets), matching_datasets=matching, removed_datasets=removed, total_size_freed_bytes=total_freed, dry_run=dry_run, errors=errors, ) if dry_run: total_size = sum( ds.raw_data_size_bytes for ds in matching if ds.raw_data_size_bytes ) msg = f"Dry run: {len(matching)} datasets ({total_size} bytes) would be removed from {db_path}." else: msg = f"Removed {len(removed)} datasets from {db_path}, freed {total_freed} bytes." log.info(msg) print(msg) return result