diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 9827f60b..92c82bc8 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -4,10 +4,38 @@ on: push: branches: [main] pull_request: + schedule: + - cron: "0 3 * * *" + workflow_dispatch: + +concurrency: + group: ci-${{ github.ref }} + cancel-in-progress: true jobs: - test: + quality-fast: + runs-on: ubuntu-latest + timeout-minutes: 20 + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-python@v5 + with: + python-version: "3.12" + cache: pip + - name: Install + run: | + python -m pip install --upgrade pip + pip install -e ".[dev,ml]" + - name: Lint + run: ruff check src tests + - name: Typecheck + run: mypy src/freshdata + - name: Test (required fast lane) + run: pytest -m "not online and not large" + + test-matrix: runs-on: ubuntu-latest + timeout-minutes: 30 strategy: fail-fast: false matrix: @@ -30,15 +58,31 @@ jobs: pip install -e ".[dev,ml]" if [ -n "${{ matrix.pandas }}" ]; then pip install "${{ matrix.pandas }}"; fi if [ -n "${{ matrix.numpy }}" ]; then pip install "${{ matrix.numpy }}"; fi - - name: Lint - run: ruff check src tests - - name: Typecheck - run: mypy src/freshdata - - name: Test - run: pytest + - name: Test (fast marker set) + run: pytest -m "not online and not large" + + nightly-online-large: + if: github.event_name == 'schedule' || github.event_name == 'workflow_dispatch' + runs-on: ubuntu-latest + timeout-minutes: 45 + continue-on-error: true + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-python@v5 + with: + python-version: "3.12" + cache: pip + - name: Install + run: | + python -m pip install --upgrade pip + pip install -e ".[dev,ml]" + - name: Nightly online/large drift checks + run: pytest -m "online or large or tier1" build: + needs: quality-fast runs-on: ubuntu-latest + timeout-minutes: 10 steps: - uses: actions/checkout@v4 - uses: actions/setup-python@v5 @@ -54,8 +98,9 @@ jobs: # Publish a self-hosted shields endpoint badge to the `badges` branch # (no third-party account needed). Runs only on main. if: github.ref == 'refs/heads/main' - needs: test + needs: quality-fast runs-on: ubuntu-latest + timeout-minutes: 15 permissions: contents: write steps: diff --git a/.github/workflows/docs.yml b/.github/workflows/docs.yml index abf298d5..5adbc30b 100644 --- a/.github/workflows/docs.yml +++ b/.github/workflows/docs.yml @@ -20,6 +20,7 @@ concurrency: jobs: build-deploy: runs-on: ubuntu-latest + timeout-minutes: 15 steps: - uses: actions/checkout@v4 with: diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 2ec81019..5a37f037 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -3,11 +3,8 @@ name: Release # Publish freshdata-cleaner to PyPI when a version tag is pushed (e.g. v0.5.0), # or manually via the Actions tab. # -# Authentication uses an API token stored in the repository secret -# PYPI_API_TOKEN (a PyPI project/account token, username __token__). `skip-existing` -# makes a re-run on an already-published version a no-op success. To switch to PyPI -# Trusted Publishing instead, configure a trusted publisher on the project and replace -# the `password:` line with `permissions: id-token: write`. +# Authentication uses PyPI Trusted Publishing (OIDC). Configure a trusted +# publisher on PyPI for this repository/workflow/environment. on: push: @@ -18,10 +15,39 @@ on: permissions: contents: read +concurrency: + group: release-${{ github.ref }} + cancel-in-progress: false + jobs: + preflight: + name: Release quality gate + runs-on: ubuntu-latest + timeout-minutes: 30 + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-python@v5 + with: + python-version: "3.12" + cache: pip + - name: Install + run: | + python -m pip install --upgrade pip + pip install -e ".[dev,ml,docs]" + - name: Lint + run: ruff check src tests + - name: Typecheck + run: mypy src/freshdata + - name: Test (release fast lane) + run: pytest -m "not online and not large" + - name: Docs strict + run: mkdocs build --strict + build: name: Build distributions + needs: preflight runs-on: ubuntu-latest + timeout-minutes: 10 steps: - uses: actions/checkout@v4 - uses: actions/setup-python@v5 @@ -41,11 +67,13 @@ jobs: name: Publish to PyPI needs: build runs-on: ubuntu-latest + timeout-minutes: 10 environment: name: pypi url: https://pypi.org/project/freshdata-cleaner/ permissions: contents: read + id-token: write steps: - uses: actions/download-artifact@v4 with: @@ -54,5 +82,4 @@ jobs: - name: Publish uses: pypa/gh-action-pypi-publish@release/v1 with: - password: ${{ secrets.PYPI_API_TOKEN }} skip-existing: true diff --git a/QUALITY_OPS.md b/QUALITY_OPS.md new file mode 100644 index 00000000..3fd720ae --- /dev/null +++ b/QUALITY_OPS.md @@ -0,0 +1,50 @@ +# Quality Operations Runbook + +This runbook defines the weekly quality cadence and the minimum metrics to track +for release readiness. + +## Weekly cadence + +- **Monday (scope/risk review)** + - Review open PR risk levels (engine changes, fixture changes, workflow changes). + - Confirm required CI lane status and top failure causes from last week. +- **Wednesday (midweek quality check)** + - Review skip counts and flaky failures. + - Verify nightly online/large job output and open drift issues. +- **Friday (merge gate)** + - Merge only PRs with passing required checks. + - Validate release readiness scorecard before tagging. + +## Scorecard metrics + +Track these weekly (rolling 4-week trend): + +- Required CI flake rate (% reruns needed to pass). +- Median required CI duration (minutes). +- `pytest` skip count in required lane. +- Nightly online/large failures (count + top 3 causes). +- Golden snapshot updates merged (count) with diff summaries attached. +- Release gate pass/fail rate. + +## Exit criteria for a release candidate + +- Required lane flake rate is near zero for the last 2 weeks. +- No unresolved deterministic regression in outlier/repair tests. +- No unexplained golden drift. +- Release workflow preflight passes on tag candidate. + +## Operational commands + +Required lane locally: + +```bash +ruff check src tests +mypy src/freshdata +pytest -m "not online and not large" +``` + +Nightly lane locally: + +```bash +pytest -m "online or large or tier1" +``` diff --git a/RELEASE.md b/RELEASE.md index 212ae60e..aa8fcf5f 100644 --- a/RELEASE.md +++ b/RELEASE.md @@ -17,7 +17,7 @@ backward-compatible fixes. ## Release checklist -1. **Green main** — `pytest`, `ruff check .`, `mypy src/freshdata`, and +1. **Green main** — `pytest -m "not online and not large"`, `ruff check .`, `mypy src/freshdata`, and `mkdocs build --strict` all pass. 2. **Bump the version** in `pyproject.toml` and `src/freshdata/__init__.py`. 3. **Update `CHANGELOG.md`** — move `Unreleased` notes under a new @@ -25,7 +25,7 @@ backward-compatible fixes. 4. **Commit & PR** — merge to `main`. 5. **Build & validate** locally (see below). 6. **Publish to TestPyPI**, smoke-test the install. -7. **Publish to PyPI** (or push the tag and let CI do it). +7. **Publish to PyPI** (push the tag and let CI do it through trusted publishing). 8. **Tag & GitHub release** — `git tag vX.Y.Z` and create the release with notes from the changelog. 9. **Verify** — `pip install freshdata-cleaner` in a clean environment imports @@ -42,10 +42,11 @@ twine check dist/* # validates metadata + long-description renderin ## Publish -### Option A — automated (recommended) +### Option A — automated (required) -Push a version tag; the `Release` workflow builds and publishes via PyPI -**Trusted Publishing** (OIDC, no stored token): +Push a version tag; the `Release` workflow runs a quality gate first +(lint, typecheck, `pytest -m "not online and not large"`, docs strict), +then builds and publishes via PyPI **Trusted Publishing** (OIDC, no stored token): ```bash git tag v0.5.0 @@ -56,7 +57,7 @@ One-time setup: add a trusted publisher at (workflow `release.yml`, environment `pypi`). -### Option B — manual with `twine` +### Option B — manual with `twine` (fallback only) ```bash # 1. TestPyPI first @@ -69,8 +70,9 @@ python -c "import freshdata as fd; print(fd.__version__)" twine upload dist/* ``` -Use a [PyPI API token](https://pypi.org/help/#apitoken) (username `__token__`). -Never commit tokens; prefer `~/.pypirc` or the `TWINE_PASSWORD` env var. +Use a [PyPI API token](https://pypi.org/help/#apitoken) (username `__token__`) only +for emergency/manual fallback releases. Never commit tokens; prefer `~/.pypirc` +or the `TWINE_PASSWORD` env var. ## Naming diff --git a/pyproject.toml b/pyproject.toml index 803fe5bb..f8deb58f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -184,6 +184,7 @@ ignore = [ [tool.ruff.lint.per-file-ignores] "tests/*" = ["PLR2004", "SIM117"] "src/freshdata/engine/model_select.py" = ["PLR0915"] +"src/freshdata/engine/outliers.py" = ["PLR0915"] "src/freshdata/adapters/polars.py" = ["PLC0415", "PLW0603"] "scripts/debug_datasets.py" = ["PLR0915"] # Enterprise layer: lazy optional imports (PLC0415) keep `import freshdata` cheap; diff --git a/scripts/debug_dataset_sweep.py b/scripts/debug_dataset_sweep.py index 2072a085..6084f423 100644 --- a/scripts/debug_dataset_sweep.py +++ b/scripts/debug_dataset_sweep.py @@ -20,8 +20,7 @@ _REPO = Path(__file__).resolve().parents[1] _TESTS = _REPO / "tests" -_LOG_PATH = _REPO.parent / ".cursor" / "debug-17298c.log" -_SESSION = "17298c" +_LOG_PATH = _REPO.parent / ".cursor" / "debug_dataset_sweep.log" if str(_TESTS) not in sys.path: sys.path.insert(0, str(_TESTS)) @@ -46,7 +45,6 @@ def _log(hypothesis_id: str, dataset: str, check: str, passed: bool, detail: str = "") -> None: _LOG_PATH.parent.mkdir(parents=True, exist_ok=True) payload = { - "sessionId": _SESSION, "hypothesisId": hypothesis_id, "location": "debug_dataset_sweep.py", "message": check, diff --git a/src/freshdata/api.py b/src/freshdata/api.py index ddeec7d2..97e9ab6c 100644 --- a/src/freshdata/api.py +++ b/src/freshdata/api.py @@ -136,6 +136,9 @@ def plan( df: pd.DataFrame, *, mode: str = "suggest", + max_patches: int | None = None, + max_cells_scanned: int | None = None, + retain_snapshots: bool = True, config: CleanConfig | None = None, **options: object, ) -> RepairPlan: @@ -147,7 +150,15 @@ def plan( deterministic representation repairs by disabling statistical engine actions. """ - return build_repair_plan(to_pandas(df), mode=mode, config=config, **options) + return build_repair_plan( + to_pandas(df), + mode=mode, + max_patches=max_patches, + max_cells_scanned=max_cells_scanned, + retain_snapshots=retain_snapshots, + config=config, + **options, + ) def repair( @@ -156,6 +167,9 @@ def repair( mode: str = "repair_safe", approved_patch_ids: set[str] | None = None, return_plan: bool = False, + max_patches: int | None = None, + max_cells_scanned: int | None = None, + retain_snapshots: bool = True, config: CleanConfig | None = None, **options: object, ) -> pd.DataFrame | tuple[pd.DataFrame, RepairPlan]: @@ -164,7 +178,15 @@ def repair( ``mode="repair_reviewed"`` applies only ``approved_patch_ids``. Other modes apply every patch proposed by the plan. """ - repair_plan = build_repair_plan(to_pandas(df), mode=mode, config=config, **options) + repair_plan = build_repair_plan( + to_pandas(df), + mode=mode, + max_patches=max_patches, + max_cells_scanned=max_cells_scanned, + retain_snapshots=retain_snapshots, + config=config, + **options, + ) approved = set(approved_patch_ids or ()) if mode == "repair_reviewed" else None repaired = from_pandas(repair_plan.apply(approved), df) if return_plan: diff --git a/src/freshdata/engine/outliers.py b/src/freshdata/engine/outliers.py index f55bc45b..e0a129dd 100644 --- a/src/freshdata/engine/outliers.py +++ b/src/freshdata/engine/outliers.py @@ -79,6 +79,10 @@ def auto_outliers(df: pd.DataFrame, config: CleanConfig, mode = config.engine_mode assert mode in ("balanced", "aggressive") rows_before = len(df) + # Detect against a stable snapshot so explicit remove actions are + # deterministic and don't depend on column processing order. + detection_df = df.copy(deep=False) + pending_remove = pd.Series(False, index=df.index) for col in list(df.columns): s = df[col] if not is_numeric_dtype(s) or is_bool_dtype(s): @@ -87,7 +91,20 @@ def auto_outliers(df: pd.DataFrame, config: CleanConfig, if len(nonnull) < _MIN_NON_NULL: continue ctx = contexts[col] - df = _handle_column(df, col, config, report, ctx=ctx, mode=mode) + df, remove_mask = _handle_column( + df, + col, + config, + report, + ctx=ctx, + mode=mode, + source_series=detection_df[col], + defer_remove=True, + ) + if remove_mask is not None and bool(remove_mask.any()): + pending_remove = pending_remove | remove_mask.reindex(df.index, fill_value=False) + if bool(pending_remove.any()): + df = df.loc[~pending_remove] removed = rows_before - len(df) if removed and removed / rows_before > _REMOVAL_WARN_SHARE: report.add_warning( @@ -137,16 +154,25 @@ def _isolation_detect(s: pd.Series, config: CleanConfig): return mask, lo, hi, "method=isolation_forest, contamination=auto" -def _handle_column(df: pd.DataFrame, col: object, config: CleanConfig, - report: CleanReport, *, ctx, mode: str) -> pd.DataFrame: - s = df[col] +def _handle_column( + df: pd.DataFrame, + col: object, + config: CleanConfig, + report: CleanReport, + *, + ctx, + mode: str, + source_series: pd.Series | None = None, + defer_remove: bool = False, +) -> tuple[pd.DataFrame, pd.Series | None]: # noqa: PLR0915 + s = source_series if source_series is not None else df[col] detected = _detect(s, config) if detected is None: - return df + return df, None mask, lo, hi, label = detected n = int(mask.sum()) if n == 0: - return df + return df, None share = n / int(s.notna().sum()) detail = f"{n} outlier(s), {100 * share:.1f}% of values ({label})" @@ -173,7 +199,7 @@ def _handle_column(df: pd.DataFrame, col: object, config: CleanConfig, f"column '{col}' has {n} extreme value(s) that were deliberately " "preserved; review them in their domain context" ) - return df + return df, None explicit = config.outlier_action not in (None, "auto") if explicit and action in ("cap", "remove") and share > _HEAVY_TAIL_SHARE: @@ -188,7 +214,7 @@ def _handle_column(df: pd.DataFrame, col: object, config: CleanConfig, risk = "low" if share <= 0.02 else "medium" if action == "cap": - df[col] = s.clip(lo, hi) + df[col] = df[col].clip(lo, hi) report.add(_STEP, f"capped {detail} to [{lo:g}, {hi:g}]", column=str(col), count=n, rationale="winsorizing keeps the rows but tames extreme " @@ -196,7 +222,11 @@ def _handle_column(df: pd.DataFrame, col: object, config: CleanConfig, risk=risk, confidence=confidence, model_id=model_id) report.outliers_handled += n elif action == "remove": - df = df.loc[~mask.fillna(False)] + if defer_remove: + remove_mask = mask.fillna(False) + else: + remove_mask = None + df = df.loc[~mask.fillna(False)] report.add(_STEP, f"removed rows with {detail}", column=str(col), count=n, rationale='outlier_action="remove" requested; rows outside ' @@ -204,12 +234,13 @@ def _handle_column(df: pd.DataFrame, col: object, config: CleanConfig, risk="medium" if share <= 0.02 else "high", confidence=confidence, model_id=model_id) report.outliers_handled += n + return df, remove_mask else: # flag base = f"{col}_outlier" if base in df.columns and df[base].dtype == bool: new_mask = mask.fillna(False).astype(bool) if df[base].equals(new_mask): - return df + return df, None flag = base else: flag = unique_flag_name(df, base) @@ -220,7 +251,7 @@ def _handle_column(df: pd.DataFrame, col: object, config: CleanConfig, "any value", risk="low", confidence=confidence, model_id=model_id) report.outliers_handled += n - return df + return df, None def _domain_sensitive(name: str) -> bool: diff --git a/src/freshdata/plan.py b/src/freshdata/plan.py index 5d3512d8..5eafcbaf 100644 --- a/src/freshdata/plan.py +++ b/src/freshdata/plan.py @@ -19,6 +19,8 @@ from .report import Action, CleanReport REPAIR_MODES = ("inspect", "suggest", "repair_safe", "repair_reviewed", "repair_aggressive") +DEFAULT_MAX_PATCHES = 20_000 +DEFAULT_MAX_CELLS_SCANNED = 2_000_000 @dataclass(frozen=True) @@ -201,6 +203,7 @@ class RepairPlan: report: CleanReport = field(default_factory=CleanReport) patches: tuple[RepairPatch, ...] = () review_items: tuple[ReviewItem, ...] = () + snapshots_retained: bool = True _before: pd.DataFrame = field(default_factory=pd.DataFrame, repr=False, compare=False) _after: pd.DataFrame = field(default_factory=pd.DataFrame, repr=False, compare=False) @@ -211,9 +214,20 @@ def patch_count(self) -> int: def apply(self, approved_patch_ids: set[str] | None = None) -> pd.DataFrame: """Apply all patches, or only the supplied approved patch identifiers.""" + if not self.snapshots_retained: + raise ValueError( + "repair plan snapshots were not retained; " + "rebuild with retain_snapshots=True to apply patches" + ) if approved_patch_ids is None: return self._after.copy(deep=True) approved = set(approved_patch_ids) + known = {patch.patch_id for patch in self.patches} + unknown = sorted(approved - known) + if unknown: + sample = ", ".join(unknown[:3]) + suffix = "..." if len(unknown) > 3 else "" + raise ValueError(f"unknown approved_patch_ids: {sample}{suffix}") out = self._before.copy(deep=True) for patch in self.patches: if patch.patch_id not in approved: @@ -230,6 +244,11 @@ def apply(self, approved_patch_ids: set[str] | None = None) -> pd.DataFrame: def rollback(self) -> pd.DataFrame: """Return the original frame captured when the plan was built.""" + if not self.snapshots_retained: + raise ValueError( + "repair plan snapshots were not retained; " + "rebuild with retain_snapshots=True to rollback" + ) return self._before.copy(deep=True) def review_queue(self) -> pd.DataFrame: @@ -420,12 +439,26 @@ def _build_patches( before: pd.DataFrame, after: pd.DataFrame, report: CleanReport, -) -> tuple[RepairPatch, ...]: + *, + max_patches: int | None = None, + max_cells_scanned: int | None = None, +) -> tuple[tuple[RepairPatch, ...], str | None]: patches: list[RepairPatch] = [] common_rows = before.index.intersection(after.index) common_cols = before.columns.intersection(after.columns) + stop_reason: str | None = None + cells_scanned = 0 + + def _can_append() -> bool: + nonlocal stop_reason + if max_patches is not None and len(patches) >= max_patches: + stop_reason = f"patch generation capped at {max_patches}" + return False + return True for row in before.index.difference(after.index): + if not _can_append(): + return tuple(patches), stop_reason action = _action_for_column(report, None) patches.append(_patch_from_action( f"p{len(patches) + 1:06d}", @@ -436,6 +469,8 @@ def _build_patches( )) for column in before.columns.difference(after.columns): + if not _can_append(): + return tuple(patches), stop_reason action = _action_for_column(report, str(column)) patches.append(_patch_from_action( f"p{len(patches) + 1:06d}", @@ -446,6 +481,8 @@ def _build_patches( )) for column in after.columns.difference(before.columns): + if not _can_append(): + return tuple(patches), stop_reason action = _action_for_column(report, str(column)) patches.append(_patch_from_action( f"p{len(patches) + 1:06d}", @@ -458,10 +495,16 @@ def _build_patches( for column in common_cols: action = _action_for_column(report, str(column)) for row in common_rows: + cells_scanned += 1 + if max_cells_scanned is not None and cells_scanned > max_cells_scanned: + stop_reason = f"cell diff scan capped at {max_cells_scanned}" + return tuple(patches), stop_reason old = before.at[row, column] new = after.at[row, column] if _values_equal(old, new): continue + if not _can_append(): + return tuple(patches), stop_reason patches.append(_patch_from_action( f"p{len(patches) + 1:06d}", "update_cell", @@ -471,7 +514,7 @@ def _build_patches( old_value=old, new_value=new, )) - return tuple(patches) + return tuple(patches), stop_reason def _build_review_items(patches: tuple[RepairPatch, ...]) -> tuple[ReviewItem, ...]: @@ -533,6 +576,9 @@ def build_repair_plan( *, mode: str = "suggest", config: CleanConfig | None = None, + max_patches: int | None = DEFAULT_MAX_PATCHES, + max_cells_scanned: int | None = DEFAULT_MAX_CELLS_SCANNED, + retain_snapshots: bool = True, **options: object, ) -> RepairPlan: """Build a repair artifact with patches, review items, and rollback data.""" @@ -547,12 +593,21 @@ def build_repair_plan( before_shape=before.shape, after_shape=before.shape, report=report, - _before=before, - _after=before, + snapshots_retained=retain_snapshots, + _before=before if retain_snapshots else pd.DataFrame(), + _after=before if retain_snapshots else pd.DataFrame(), ) run_cfg = merge_options(cfg, verbose=False, preserve_original=True) after, report = run_pipeline(before, run_cfg) - patches = _build_patches(before, after, report) + patches, stop_reason = _build_patches( + before, + after, + report, + max_patches=max_patches, + max_cells_scanned=max_cells_scanned, + ) + if stop_reason: + report.add_warning(f"repair patch diff truncated: {stop_reason}") review_items = _build_review_items(patches) return RepairPlan( config=cfg, @@ -563,8 +618,9 @@ def build_repair_plan( report=report, patches=patches, review_items=review_items, - _before=before, - _after=after.copy(deep=True), + snapshots_retained=retain_snapshots, + _before=before if retain_snapshots else pd.DataFrame(), + _after=after.copy(deep=True) if retain_snapshots else pd.DataFrame(), ) diff --git a/tests/conftest.py b/tests/conftest.py index 29222a20..791b7e74 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -9,6 +9,12 @@ def pytest_addoption(parser): default=False, help="Rewrite golden report snapshots in tests/fixtures/golden/ and online/golden/", ) + parser.addoption( + "--require-golden-diff", + action="store_true", + default=False, + help="Fail golden updates unless a machine-readable diff summary is written.", + ) @pytest.fixture @@ -16,6 +22,11 @@ def update_golden(request): return request.config.getoption("--update-golden") +@pytest.fixture +def require_golden_diff(request): + return request.config.getoption("--require-golden-diff") + + @pytest.fixture def messy() -> pd.DataFrame: """A kitchen-sink frame exercising every default cleaning step.""" diff --git a/tests/expectations.py b/tests/expectations.py index d4881003..3a74d329 100644 --- a/tests/expectations.py +++ b/tests/expectations.py @@ -4,6 +4,7 @@ import hashlib import json +import os import sys import time import urllib.request @@ -63,7 +64,7 @@ def load_online_manifest() -> dict: def load_fixture(name: str) -> pd.DataFrame: path = FIXTURES_DIR / f"{name}.csv" if not path.exists(): - pytest.skip(f"fixture {name} not found") + pytest.fail(f"fixture {name} not found at {path}") return pd.read_csv(path) @@ -94,12 +95,15 @@ def _fetch_online_live(name: str) -> pd.DataFrame: def load_online_fixture(name: str, *, live: bool = False) -> pd.DataFrame: """Load cached online slice; optionally fetch live when *live* is True.""" cache_path = ONLINE_CACHE_DIR / f"{name}.csv" + strict_online = os.environ.get("FRESHDATA_STRICT_ONLINE_FIXTURES", "0") == "1" if cache_path.exists() and not live: df = pd.read_csv(cache_path) df.columns = [str(c) for c in df.columns] return df if live: return _fetch_online_live(name) + if strict_online: + pytest.fail(f"online cache missing for {name!r}; run scripts/fetch_online_fixtures.py") pytest.skip(f"online cache missing for {name!r}; run scripts/fetch_online_fixtures.py") diff --git a/tests/golden_util.py b/tests/golden_util.py index 26c42b73..33bb07cd 100644 --- a/tests/golden_util.py +++ b/tests/golden_util.py @@ -12,6 +12,7 @@ GOLDEN_DIR = FIXTURES_DIR / "golden" ONLINE_GOLDEN_DIR = ONLINE_DIR / "golden" +GOLDEN_DIFF_SUMMARY_PATH = FIXTURES_DIR / "golden_diff_summary.jsonl" def normalize_report(report: fd.CleanReport) -> dict[str, Any]: @@ -55,5 +56,21 @@ def write_golden( ) -> Path: path = golden_path(fixture_name, strategy, online=online) path.parent.mkdir(parents=True, exist_ok=True) - path.write_text(json.dumps(normalize_report(report), indent=2, sort_keys=True) + "\n") + normalized = normalize_report(report) + previous = json.loads(path.read_text()) if path.exists() else None + path.write_text(json.dumps(normalized, indent=2, sort_keys=True) + "\n") + diff_summary = { + "fixture": fixture_name, + "strategy": strategy, + "online": online, + "created": previous is None, + "changed": previous != normalized, + "previous_action_count": ( + len(previous.get("actions", [])) if isinstance(previous, dict) else 0 + ), + "new_action_count": len(normalized.get("actions", [])), + } + GOLDEN_DIFF_SUMMARY_PATH.parent.mkdir(parents=True, exist_ok=True) + with GOLDEN_DIFF_SUMMARY_PATH.open("a", encoding="utf-8") as fh: + fh.write(json.dumps(diff_summary, sort_keys=True) + "\n") return path diff --git a/tests/test_compare_matrix.py b/tests/test_compare_matrix.py index 1d667f33..d544577e 100644 --- a/tests/test_compare_matrix.py +++ b/tests/test_compare_matrix.py @@ -21,7 +21,9 @@ def test_compare_clean_all_strategies(fixture_name): assert set(table["strategy"]) >= {"conservative", "balanced", "aggressive"} assert (table["duration_seconds"] >= 0).all() conservative = table.loc[table["strategy"] == "conservative"].iloc[0] - assert conservative["missing_after"] >= conservative["missing_before"] - 1 + assert conservative["rows_before"] == len(df) + assert conservative["rows_after"] > 0 + assert isinstance(json.loads(conservative["primary_models"]), dict) @pytest.mark.parametrize("fixture_name", ["aqi_sample", "large_panel", "wide_sparse"]) diff --git a/tests/test_config_tuple_fixtures.py b/tests/test_config_tuple_fixtures.py index f325b6cf..78750aea 100644 --- a/tests/test_config_tuple_fixtures.py +++ b/tests/test_config_tuple_fixtures.py @@ -28,9 +28,7 @@ def _parity(clean_a, clean_b) -> None: def test_id_columns_string_matches_tuple(fixture_name): df = load_fixture(fixture_name) id_cols = [c for c in df.columns if c.lower().endswith("_id") or c.lower() == "id"] - if not id_cols: - pytest.skip("no id-like column") - col = str(id_cols[0]) + col = str(id_cols[0] if id_cols else df.columns[0]) kw = {**ISOLATE, "id_columns": col, "outlier_action": "cap"} _parity(fd.clean(df, **kw), fd.clean(df, **{**kw, "id_columns": (col,)})) @@ -40,9 +38,7 @@ def test_preserve_columns_string_matches_tuple(fixture_name): df = load_fixture(fixture_name) exp = load_expectations(fixture_name).get("balanced", {}) cols = exp.get("columns_never_imputed") or exp.get("columns_preserved_as_text") or [] - if not cols: - pytest.skip("no preserve expectations") - col = resolve_column(df, cols[0]) + col = resolve_column(df, cols[0]) if cols else str(df.columns[-1]) kw = {**ISOLATE, "preserve_columns": col, "outlier_action": "cap"} _parity(fd.clean(df, **kw), fd.clean(df, **{**kw, "preserve_columns": (col,)})) diff --git a/tests/test_engine_outliers.py b/tests/test_engine_outliers.py index dfd23ab7..add34692 100644 --- a/tests/test_engine_outliers.py +++ b/tests/test_engine_outliers.py @@ -53,6 +53,19 @@ def test_remove_action_drops_rows(): assert "removed" in action.description +def test_remove_action_is_column_order_invariant(): + vals_a = [1.0] * 100 + [1000.0] + vals_b = [1.0] * 100 + [2000.0] + base = pd.DataFrame({"a": vals_a, "b": vals_b}) + swapped = base[["b", "a"]] + + out_base = fd.clean(base, outlier_action="remove", **QUIET) + out_swapped = fd.clean(swapped, outlier_action="remove", **QUIET) + + assert len(out_base) == len(out_swapped) + assert out_base.reset_index(drop=True).equals(out_swapped[["a", "b"]].reset_index(drop=True)) + + def test_flag_action_keeps_data(): df = pd.DataFrame({"v": normal_with_spike()}) out = fd.clean(df, outlier_action="flag", **QUIET) diff --git a/tests/test_expectations_preflight.py b/tests/test_expectations_preflight.py new file mode 100644 index 00000000..f2247f1e --- /dev/null +++ b/tests/test_expectations_preflight.py @@ -0,0 +1,19 @@ +"""Fixture preflight behavior for expectation helpers.""" + +from __future__ import annotations + +import pytest +from _pytest.outcomes import Failed + +from expectations import load_fixture, load_online_fixture + + +def test_load_fixture_fails_when_local_fixture_missing(): + with pytest.raises(Failed, match="fixture definitely_missing_fixture not found"): + load_fixture("definitely_missing_fixture") + + +def test_load_online_fixture_strict_mode_fails_when_cache_missing(monkeypatch): + monkeypatch.setenv("FRESHDATA_STRICT_ONLINE_FIXTURES", "1") + with pytest.raises(Failed, match="online cache missing"): + load_online_fixture("definitely_missing_online_fixture") diff --git a/tests/test_golden.py b/tests/test_golden.py index 69eb150c..6e390199 100644 --- a/tests/test_golden.py +++ b/tests/test_golden.py @@ -8,22 +8,25 @@ from __future__ import annotations import json +from pathlib import Path import pytest import freshdata as fd from expectations import ALL_FIXTURES, load_fixture -from golden_util import load_golden, normalize_report, write_golden +from golden_util import GOLDEN_DIFF_SUMMARY_PATH, load_golden, normalize_report, write_golden @pytest.mark.parametrize("fixture_name", ALL_FIXTURES) -def test_balanced_report_golden_snapshot(fixture_name, update_golden): +def test_balanced_report_golden_snapshot(fixture_name, update_golden, require_golden_diff): df = load_fixture(fixture_name) _, report = fd.clean(df, return_report=True, verbose=False) actual = normalize_report(report) if update_golden: path = write_golden(fixture_name, report, strategy="balanced") + if require_golden_diff: + assert Path(GOLDEN_DIFF_SUMMARY_PATH).exists(), "golden diff summary was not produced" pytest.skip(f"updated golden snapshot: {path}") expected = load_golden(fixture_name, strategy="balanced") diff --git a/tests/test_plan.py b/tests/test_plan.py index 5bae2386..df2341b0 100644 --- a/tests/test_plan.py +++ b/tests/test_plan.py @@ -100,6 +100,14 @@ def test_repair_reviewed_without_approvals_preserves_input(): assert_frame_equal(repaired, df) +def test_repair_reviewed_unknown_patch_id_raises(): + df = pd.DataFrame({"name": [" Ann "]}) + plan = fd.plan(df, mode="repair_safe") + assert plan.patches + with pytest.raises(ValueError, match="unknown approved_patch_ids"): + plan.apply({"does-not-exist"}) + + def test_repair_plan_review_queue_and_dbt_export_shapes(): df = pd.DataFrame({"value": [1, 2, 1000, None]}) plan = fd.plan(df, mode="repair_aggressive") @@ -116,3 +124,28 @@ def test_repair_plan_inspect_mode_does_not_propose_patches(): assert plan.patch_count == 0 assert_frame_equal(plan.apply(), df) + + +def test_repair_plan_patch_generation_limits_warn_and_truncate(): + df = pd.DataFrame({"name": [" Ann ", " Bob ", " Cat "]}) + repaired, plan = fd.repair( + df, + mode="repair_safe", + max_patches=1, + return_plan=True, + strategy="conservative", + ) + assert plan.patch_count == 1 + assert isinstance(repaired, pd.DataFrame) + assert any("repair patch diff truncated" in w for w in plan.report.warnings) + + +def test_repair_plan_without_snapshots_blocks_apply_and_rollback(): + df = pd.DataFrame({"name": [" Ann "]}) + plan = fd.plan(df, mode="repair_safe", retain_snapshots=False) + assert plan.patch_count >= 0 + assert plan.snapshots_retained is False + with pytest.raises(ValueError, match="snapshots were not retained"): + plan.apply() + with pytest.raises(ValueError, match="snapshots were not retained"): + plan.rollback() diff --git a/tests/test_realworld.py b/tests/test_realworld.py index d4cbda1e..f1191a88 100644 --- a/tests/test_realworld.py +++ b/tests/test_realworld.py @@ -16,8 +16,11 @@ def test_clean_produces_valid_output(fixture_name, strategy): df = load_fixture(fixture_name) out, report, duration = clean_with_timing(df, strategy=strategy) - assert len(out) >= 0 + assert report.rows_after == len(out) + assert report.cols_after == out.shape[1] assert report.rows_before == len(df) + assert out.columns.is_unique + assert report.duration_seconds >= 0 assert_expectations(fixture_name, strategy, df, out, report, duration=duration) diff --git a/tests/test_report.py b/tests/test_report.py index 9736020f..e2b04a92 100644 --- a/tests/test_report.py +++ b/tests/test_report.py @@ -23,8 +23,10 @@ def test_report_is_iterable_and_sized(messy): def test_to_dict_is_json_serializable(messy): _, report = fd.clean(messy, report=True) - payload = json.dumps(report.to_dict()) + payload_dict = report.to_dict() + payload = json.dumps(payload_dict) assert "drop_duplicates" in payload + assert len(payload_dict["actions"]) == len(report) def test_to_frame(messy):