From f6104d1f5697a3cef4211e96839bee6eddee8db2 Mon Sep 17 00:00:00 2001 From: Kevin Costner <120246174+kevincostner17@users.noreply.github.com> Date: Tue, 15 Sep 2026 21:55:41 +0530 Subject: [PATCH 1/5] fix(dbt): keep dbt-gate stdout to a single JSON document evaluate_trust_gate called fd.clean with clean_config=None, so CleanConfig's verbose=True default printed the clean summary and its warnings to stdout ahead of dbt-gate's JSON summary, and 'dbt-gate | jq' failed. The shared gate (dbt, Airflow, Dagster) now cleans quietly when no config is given and logs the clean warnings through the freshdata.integrations logger; an explicit CleanConfig is still used as given. dbt-gate also routes anything printed while gating to stderr, so stdout carries only the summary. --- docs/integrations.md | 4 +++- src/freshdata/integrations/_core.py | 12 +++++++++++- src/freshdata/integrations/dbt/cli.py | 18 +++++++++++------- tests/test_integrations/test_core.py | 20 ++++++++++++++++++++ tests/test_integrations/test_dbt.py | 26 ++++++++++++++++++++++++++ 5 files changed, 71 insertions(+), 9 deletions(-) diff --git a/docs/integrations.md b/docs/integrations.md index 56c6de2c..ff563ffd 100644 --- a/docs/integrations.md +++ b/docs/integrations.md @@ -100,7 +100,9 @@ gates it, and exits non-zero (with `--fail`) if any model is below the threshold Ephemeral and disabled models are not read; they are listed under `"skipped"`. If no model is gated at all (an empty manifest, or only ephemeral models), `all_passed` is `false` and `--fail` exits 1. A file that is not a dbt manifest (for example -`run_results.json`) is reported as a one-line error with exit 1. For +`run_results.json`) is reported as a one-line error with exit 1. Stdout carries +only the JSON summary, so it can be piped straight into `jq` or `json.loads`; +cleaning warnings and other messages go to stderr. For a single model — or to write per-model `_audit.json` files — use `FreshDataDbtTransform`: diff --git a/src/freshdata/integrations/_core.py b/src/freshdata/integrations/_core.py index db22e8d9..47a59f60 100644 --- a/src/freshdata/integrations/_core.py +++ b/src/freshdata/integrations/_core.py @@ -200,7 +200,17 @@ def evaluate_trust_gate( """ on_low_score = validate_on_low_score(on_low_score) row_count_in = len(df) - cleaned, report = fd.clean(df, config=clean_config, report=True) + # The gate reports through TrustGateResult and logging, never stdout: callers + # such as the ``dbt-gate`` CLI promise machine-readable stdout. Without an + # explicit config, clean quietly; a caller's own CleanConfig is used as given. + if clean_config is None: + cleaned, report = fd.clean(df, report=True, verbose=False) + else: + cleaned, report = fd.clean(df, config=clean_config, report=True) + if clean_config is None or not clean_config.verbose: + # Quiet cleans still surface their warnings, via logging (stderr by default). + for warning in report.warnings: + logger.warning("freshdata: %s", warning) trust = compute_trust_score(cleaned) score = float(trust.overall) high_risk = sum(1 for action in report.actions if getattr(action, "risk", None) == "high") diff --git a/src/freshdata/integrations/dbt/cli.py b/src/freshdata/integrations/dbt/cli.py index 4e401870..ac4d21a5 100644 --- a/src/freshdata/integrations/dbt/cli.py +++ b/src/freshdata/integrations/dbt/cli.py @@ -11,6 +11,7 @@ from __future__ import annotations import argparse +import contextlib import json import sys @@ -61,13 +62,16 @@ def main(argv: list[str] | None = None) -> int: """Entry point for the ``dbt-gate`` script. Returns a process exit code.""" args = _build_parser().parse_args(argv) try: - summary = gate_manifest( - args.manifest, - conn_str=args.conn, - trust_score_threshold=args.threshold, - on_low_score=args.on_low_score, - output_dir=args.output_dir, - ) + # stdout carries exactly one JSON document. Anything printed while gating + # (a verbose clean, a third-party library) goes to stderr instead. + with contextlib.redirect_stdout(sys.stderr): + summary = gate_manifest( + args.manifest, + conn_str=args.conn, + trust_score_threshold=args.threshold, + on_low_score=args.on_low_score, + output_dir=args.output_dir, + ) except (OSError, ValueError) as exc: # A wrong, unreadable or malformed manifest (missing file, a directory, # invalid JSON, not a dbt manifest) is routine CLI misuse, not a crash: diff --git a/tests/test_integrations/test_core.py b/tests/test_integrations/test_core.py index 39146990..d518b8c9 100644 --- a/tests/test_integrations/test_core.py +++ b/tests/test_integrations/test_core.py @@ -2,10 +2,12 @@ from __future__ import annotations +import logging import sys import pytest +from freshdata import CleanConfig from freshdata.integrations import TrustGateResult, evaluate_trust_gate @@ -96,3 +98,21 @@ def test_no_report_dict_without_publish(sample_df): _, result = evaluate_trust_gate(sample_df, trust_score_threshold=0.0) assert result.report_dict is None assert "report" not in result.to_dict() + + +def test_gate_without_config_keeps_stdout_clean_and_logs_warnings(sample_df, capsys, caplog): + """The gate must not print: ``dbt-gate`` promises JSON-only stdout.""" + caplog.set_level(logging.WARNING, logger="freshdata.integrations") + _, result = evaluate_trust_gate( + sample_df, trust_score_threshold=0.0, publish_full_report=True + ) + assert capsys.readouterr().out == "" + warnings = result.report_dict["clean_report"]["warnings"] + assert warnings # sample_df's duplicate row exceeds the duplicate threshold + for warning in warnings: + assert f"freshdata: {warning}" in caplog.messages + + +def test_gate_respects_explicit_verbose_config(sample_df, capsys): + evaluate_trust_gate(sample_df, clean_config=CleanConfig(verbose=True)) + assert "freshdata: rows" in capsys.readouterr().out diff --git a/tests/test_integrations/test_dbt.py b/tests/test_integrations/test_dbt.py index 4fd81507..3287e803 100644 --- a/tests/test_integrations/test_dbt.py +++ b/tests/test_integrations/test_dbt.py @@ -121,6 +121,32 @@ def test_cli_pass_and_fail_exit_codes(warehouse, tmp_path, capsys): assert rc == 1 +def test_cli_stdout_is_exactly_one_json_document(warehouse, tmp_path, capsys): + # sample_df has a duplicate row, so cleaning raises a warning; before the fix it + # was printed to stdout ahead of the summary and broke `dbt-gate | jq`. + manifest = str(_manifest(tmp_path)) + assert main(["--manifest", manifest, "--conn", warehouse, "--threshold", "0"]) == 0 + out = capsys.readouterr().out + summary = json.loads(out) + assert summary["models_processed"] == 1 + assert out.lstrip().startswith("{") + + +def test_cli_prints_during_gating_go_to_stderr(tmp_path, capsys, monkeypatch): + import freshdata.integrations.dbt.cli as dbt_cli + + def noisy_gate(*args, **kwargs): + print("chatty library output") + return {"models": [], "skipped": [], "models_processed": 1, "failed_models": 0, + "all_passed": True} + + monkeypatch.setattr(dbt_cli, "gate_manifest", noisy_gate) + assert main(["--manifest", str(tmp_path / "manifest.json")]) == 0 + captured = capsys.readouterr() + assert json.loads(captured.out)["all_passed"] is True + assert "chatty library output" in captured.err + + def test_cli_missing_manifest_prints_one_line_error(capsys): code = main(["--manifest", "definitely_not_here.json"]) assert code == 1 From 9f2b4a77fd234578d9ded931919adfa85d62aa42 Mon Sep 17 00:00:00 2001 From: Kevin Costner <120246174+kevincostner17@users.noreply.github.com> Date: Tue, 15 Sep 2026 22:05:18 +0530 Subject: [PATCH 2/5] fix(cli): reject unknown keys in freshdata clean --config _build_enterprise read only six scalar keys plus masking, semantic and clustering from the 'enterprise' section, and the loader never looked at top-level section names. Typos such as enable_maskin or fail_under_trus, a misspelled 'enterprize:' section, and real EnterpriseConfig fields such as enable_privacy_detection or privacy were all silently ignored with exit 0. Unknown sections and keys, including keys inside nested objects, now fail with a one-line did-you-mean error (exit 1), worded like the 'clean' section's error. Every EnterpriseConfig field is either applied or rejected: enable_privacy_detection, enable_entity_resolution, trust_weights, lineage, privacy, k_anonymity and entity_resolution are now built into the config; enable_contracts, drift and anonymization are rejected with the reason, because the clean command never runs contract checks or applies them. Boolean, actor and fail_under_trust values are type-checked. The accepted keys are documented in docs/feature-overview.md. --- docs/feature-overview.md | 39 +++++++ src/freshdata/enterprise/cli.py | 164 ++++++++++++++++++++++++++--- tests/test_cli_malformed_inputs.py | 143 +++++++++++++++++++++++++ 3 files changed, 334 insertions(+), 12 deletions(-) diff --git a/docs/feature-overview.md b/docs/feature-overview.md index 4506542c..a3b9aa62 100644 --- a/docs/feature-overview.md +++ b/docs/feature-overview.md @@ -48,6 +48,45 @@ print(result.quality.to_markdown()) assert result.passed_gate ``` +### `freshdata clean --config` files + +`--config` takes a JSON or YAML object with two optional sections, `clean` and +`enterprise`. An unknown section or key, including a typo, stops the run before +any data is read: a one-line error with a "did you mean" hint, exit 1. + +```yaml +clean: + strategy: balanced +enterprise: + fail_under_trust: 80 + masking: + - {name: pii, columns: [email], strategy: hash} + enable_privacy_detection: true + privacy: {min_score: 0.6} +``` + +`clean` accepts any `CleanConfig` option. `enterprise` accepts these keys: + +| Key | Value | Builds | +|---|---|---| +| `actor` | string or null | `EnterpriseConfig.actor` | +| `fail_under_trust` | number from 0 to 100, or null | the trust gate; `--fail-under-trust` overrides it | +| `enable_masking`, `enable_clustering`, `enable_validation`, `enable_lineage`, `enable_privacy_detection`, `enable_entity_resolution` | `true` or `false` | the toggle of the same name | +| `masking` | list of objects | one `MaskingRule` each; `--mask` adds more | +| `semantic` | list of objects | one `SemanticValidatorConfig` each | +| `clustering` | object | `ClusterConfig`; `--cluster` replaces it | +| `trust_weights` | object | `TrustScoreWeights` | +| `lineage` | object | `LineageConfig` | +| `privacy` | object | `PIIDetectionConfig`, applied when `enable_privacy_detection` is true | +| `k_anonymity` | object | `KAnonymityConfig` | +| `entity_resolution` | object, with `blocking_rules` and `comparisons` as lists of objects | `EntityResolutionConfig` (with `BlockingRule` and `ComparisonLevel`), applied when `enable_entity_resolution` is true | + +Nested objects take the field names of the class they build, and unknown names +are rejected the same way. Three `EnterpriseConfig` fields are rejected with an +explanation: `enable_contracts` and `drift` (`freshdata clean` takes no baseline or +data contract to check against) and `anonymization` (no pipeline applies it; use +`masking`, or `privacy` with `enable_privacy_detection`). + ## Compliance reports The `freshdata.compliance` subpackage turns a `CleanReport` into a regulatory diff --git a/src/freshdata/enterprise/cli.py b/src/freshdata/enterprise/cli.py index 1899e48e..9521cfcf 100644 --- a/src/freshdata/enterprise/cli.py +++ b/src/freshdata/enterprise/cli.py @@ -14,8 +14,11 @@ from __future__ import annotations import argparse +import dataclasses +import difflib import json import sys +from collections.abc import Collection from pathlib import Path from typing import Any @@ -27,7 +30,19 @@ from ..context import PolicyError from ..insight import insight_report, trust_gate_report from ..profile import build_profile -from .config import ClusterConfig, EnterpriseConfig, MaskingRule, SemanticValidatorConfig +from .config import ( + BlockingRule, + ClusterConfig, + ComparisonLevel, + EnterpriseConfig, + EntityResolutionConfig, + KAnonymityConfig, + LineageConfig, + MaskingRule, + PIIDetectionConfig, + SemanticValidatorConfig, + TrustScoreWeights, +) from .interface import clean_enterprise from .metrics import compute_trust_score @@ -173,20 +188,144 @@ def _config_section(data: dict[str, Any], key: str, path: str) -> dict[str, Any] return section -def _build_enterprise(spec: dict[str, Any]) -> EnterpriseConfig: - masking = tuple(MaskingRule(**rule) for rule in spec.get("masking", [])) - semantic = tuple(SemanticValidatorConfig(**val) for val in spec.get("semantic", [])) - clustering = ClusterConfig(**spec["clustering"]) if spec.get("clustering") else None - scalar_keys = ( +#: Top-level sections a ``freshdata clean --config`` file may contain. +_CONFIG_SECTIONS = ("clean", "enterprise") + +#: ``enterprise`` keys taken as plain booleans. +_ENTERPRISE_BOOL_KEYS = ( + "enable_masking", + "enable_clustering", + "enable_validation", + "enable_lineage", + "enable_privacy_detection", + "enable_entity_resolution", +) +#: ``enterprise`` keys holding one object, built into the named dataclass. +_ENTERPRISE_OBJECT_KEYS: dict[str, Any] = { + "clustering": ClusterConfig, + "trust_weights": TrustScoreWeights, + "lineage": LineageConfig, + "privacy": PIIDetectionConfig, + "k_anonymity": KAnonymityConfig, + "entity_resolution": EntityResolutionConfig, +} +#: ``enterprise`` keys holding a list of objects, each built into the named dataclass. +_ENTERPRISE_LIST_KEYS: dict[str, Any] = { + "masking": MaskingRule, + "semantic": SemanticValidatorConfig, +} +#: List-of-object fields inside a nested config, built the same way. +_NESTED_LIST_FIELDS: dict[Any, dict[str, Any]] = { + EntityResolutionConfig: {"blocking_rules": BlockingRule, "comparisons": ComparisonLevel}, +} +#: Real ``EnterpriseConfig`` fields ``freshdata clean`` cannot honour, with the reason. +_ENTERPRISE_UNSUPPORTED_KEYS = { + "enable_contracts": ( + "contract checks need a baseline or data contract, which 'freshdata clean' does not take" + ), + "drift": ( + "drift thresholds only apply to contract checks, which 'freshdata clean' does not run" + ), + "anonymization": ( + "no pipeline applies it; use 'masking', or 'privacy' with 'enable_privacy_detection'" + ), +} +_ENTERPRISE_KEYS = frozenset( + ( "actor", - "enable_masking", - "enable_clustering", - "enable_validation", - "enable_lineage", "fail_under_trust", + *_ENTERPRISE_BOOL_KEYS, + *_ENTERPRISE_OBJECT_KEYS, + *_ENTERPRISE_LIST_KEYS, + ) +) + + +def _check_keys( + spec: dict[Any, Any], valid: Collection[str], what: str, *, suggest: Collection[str] = () +) -> None: + """Raise :class:`TypeError` naming keys of *spec* outside *valid*. + + Worded like the ``clean`` section's error from :func:`freshdata.config.merge_options`: + each unknown key gets a "did you mean" hint drawn from *valid* plus *suggest*. + """ + unknown = sorted({str(key) for key in spec} - set(valid)) + if not unknown: + return + pool = sorted({*valid, *suggest}) + hints = [] + for name in unknown: + match = difflib.get_close_matches(name, pool, n=1) + hints.append(f"{name!r}" + (f" (did you mean {match[0]!r}?)" if match else "")) + raise TypeError( + f"unknown {what}(s): {', '.join(hints)}. Valid {what}s: {', '.join(sorted(valid))}" ) - kwargs = {key: spec[key] for key in scalar_keys if key in spec} - return EnterpriseConfig(masking=masking, semantic=semantic, clustering=clustering, **kwargs) + + +def _build_dataclass(cls: Any, spec: Any, where: str) -> Any: + """Build config dataclass *cls* from the mapping *spec* found at *where*.""" + if not isinstance(spec, dict): + raise TypeError(f"'{where}' must be an object, got {type(spec).__name__}") + _check_keys(spec, [f.name for f in dataclasses.fields(cls)], f"'{where}' key") + kwargs = dict(spec) + for key, item_cls in _NESTED_LIST_FIELDS.get(cls, {}).items(): + if key in kwargs: + kwargs[key] = _build_list(item_cls, kwargs[key], f"{where}.{key}") + return cls(**kwargs) + + +def _build_list(cls: Any, items: Any, where: str) -> tuple[Any, ...]: + """Build a tuple of *cls* from the list of mappings *items* found at *where*.""" + if not isinstance(items, list): + raise TypeError(f"'{where}' must be a list, got {type(items).__name__}") + return tuple(_build_dataclass(cls, item, f"{where}[{i}]") for i, item in enumerate(items)) + + +def _check_config_sections(data: dict[str, Any], path: str) -> None: + """Reject unknown top-level sections (e.g. a misspelled ``enterprize:``).""" + try: + _check_keys(data, _CONFIG_SECTIONS, "section") + except TypeError as exc: + raise ValueError(f"invalid config file {path}: {exc}") from exc + + +def _build_enterprise(spec: dict[str, Any]) -> EnterpriseConfig: + """Build an :class:`EnterpriseConfig` from a ``--config`` file's ``enterprise`` section. + + Every key must be one ``freshdata clean`` applies. Typos, and real + ``EnterpriseConfig`` fields this command cannot apply, raise :class:`TypeError` + instead of being silently ignored. + """ + for key, reason in _ENTERPRISE_UNSUPPORTED_KEYS.items(): + if key in spec: + raise TypeError(f"{key!r} is not supported in --config: {reason}") + _check_keys(spec, _ENTERPRISE_KEYS, "key", suggest=_ENTERPRISE_UNSUPPORTED_KEYS) + + kwargs: dict[str, Any] = {} + for key in _ENTERPRISE_BOOL_KEYS: + if key in spec: + if not isinstance(spec[key], bool): + raise TypeError(f"'{key}' must be true or false, got {spec[key]!r}") + kwargs[key] = spec[key] + if "actor" in spec: + if spec["actor"] is not None and not isinstance(spec["actor"], str): + raise TypeError(f"'actor' must be a string or null, got {spec['actor']!r}") + kwargs["actor"] = spec["actor"] + if "fail_under_trust" in spec: + value = spec["fail_under_trust"] + if value is not None and (isinstance(value, bool) or not isinstance(value, (int, float))): + raise TypeError(f"'fail_under_trust' must be a number or null, got {value!r}") + kwargs["fail_under_trust"] = value + for key, cls in _ENTERPRISE_LIST_KEYS.items(): + if key in spec: + kwargs[key] = _build_list(cls, spec[key], key) + for key, cls in _ENTERPRISE_OBJECT_KEYS.items(): + value = spec.get(key) + # An empty ``clustering`` object has always meant "no clustering". + if value is None or (key == "clustering" and not value): + continue + kwargs[key] = _build_dataclass(cls, value, key) + return EnterpriseConfig(**kwargs) def _load_profile_arg(path: str, *, quiet: bool = False) -> tuple[Any, int]: @@ -228,6 +367,7 @@ def cmd_clean(args: argparse.Namespace) -> int: ec = EnterpriseConfig() if args.config: data = _load_config_file(args.config) + _check_config_sections(data, args.config) file_clean = _config_section(data, "clean", args.config) try: ec = _build_enterprise(_config_section(data, "enterprise", args.config)) diff --git a/tests/test_cli_malformed_inputs.py b/tests/test_cli_malformed_inputs.py index af8b1cd6..8e62a4a2 100644 --- a/tests/test_cli_malformed_inputs.py +++ b/tests/test_cli_malformed_inputs.py @@ -8,6 +8,7 @@ from __future__ import annotations +import dataclasses import io import json import sys @@ -16,6 +17,7 @@ import pytest from freshdata.enterprise import cli +from freshdata.enterprise.config import EnterpriseConfig from freshdata.validation_suite import ValidationSuite @@ -123,6 +125,147 @@ def test_clean_valid_config_sections_still_work(src, tmp_path): assert "a@b.com" not in (tmp_path / "out.csv").read_text() +# --------------------------------------------------------------------------- # +# Unknown / unsupported 'enterprise' keys and top-level sections # +# --------------------------------------------------------------------------- # +def _write_cfg(tmp_path, payload): + cfg = tmp_path / "cfg.json" + cfg.write_text(json.dumps(payload)) + return cfg + + +def test_clean_unknown_enterprise_keys_error_with_did_you_mean(src, tmp_path, capsys): + cfg = _write_cfg(tmp_path, {"enterprise": {"enable_maskin": True, "fail_under_trus": 99}}) + assert _clean_with_config(src, cfg, tmp_path) == 1 + err = _assert_one_line_error( + capsys, + "cfg.json", + "'enable_maskin' (did you mean 'enable_masking'?)", + "'fail_under_trus' (did you mean 'fail_under_trust'?)", + ) + assert len(err.strip().splitlines()) == 1 + assert not (tmp_path / "out.csv").exists() + + +def test_clean_unknown_top_level_section_errors_with_did_you_mean(src, tmp_path, capsys): + cfg = _write_cfg(tmp_path, {"enterprize": {"enable_masking": True}}) + assert _clean_with_config(src, cfg, tmp_path) == 1 + err = _assert_one_line_error( + capsys, "cfg.json", "'enterprize' (did you mean 'enterprise'?)", "clean, enterprise" + ) + assert len(err.strip().splitlines()) == 1 + assert not (tmp_path / "out.csv").exists() + + +@pytest.mark.parametrize( + "payload, needle", + [ + ({"privacy": {"min_scor": 0.5}}, "'min_scor' (did you mean 'min_score'?)"), + ({"clustering": {"colums": ["email"]}}, "'colums' (did you mean 'columns'?)"), + ( + {"entity_resolution": {"comparisons": [{"colum": "email"}]}}, + "'colum' (did you mean 'column'?)", + ), + ], +) +def test_clean_unknown_nested_enterprise_key_is_one_line_error( + src, tmp_path, capsys, payload, needle +): + cfg = _write_cfg(tmp_path, {"enterprise": payload}) + assert _clean_with_config(src, cfg, tmp_path) == 1 + _assert_one_line_error(capsys, "cfg.json", needle) + + +@pytest.mark.parametrize( + "key, value", [("enable_contracts", True), ("drift", {}), ("anonymization", [])] +) +def test_clean_unsupported_enterprise_field_is_rejected(src, tmp_path, capsys, key, value): + cfg = _write_cfg(tmp_path, {"enterprise": {key: value}}) + assert _clean_with_config(src, cfg, tmp_path) == 1 + _assert_one_line_error(capsys, "cfg.json", f"'{key}' is not supported in --config") + + +@pytest.mark.parametrize( + "payload, needle", + [ + ({"enable_masking": "false"}, "'enable_masking' must be true or false"), + ({"fail_under_trust": "80"}, "'fail_under_trust' must be a number or null"), + ({"actor": 7}, "'actor' must be a string or null"), + ({"privacy": [1]}, "'privacy' must be an object"), + ], +) +def test_clean_wrongly_typed_enterprise_value_is_one_line_error( + src, tmp_path, capsys, payload, needle +): + cfg = _write_cfg(tmp_path, {"enterprise": payload}) + assert _clean_with_config(src, cfg, tmp_path) == 1 + _assert_one_line_error(capsys, "cfg.json", needle) + + +def test_every_enterprise_config_field_is_accepted_or_rejected(): + """A new EnterpriseConfig field must be wired into --config or rejected explicitly.""" + fields = {f.name for f in dataclasses.fields(EnterpriseConfig)} + handled = set(cli._ENTERPRISE_KEYS) | set(cli._ENTERPRISE_UNSUPPORTED_KEYS) + assert fields == handled + assert not set(cli._ENTERPRISE_KEYS) & set(cli._ENTERPRISE_UNSUPPORTED_KEYS) + + +def _clean_loud(src, cfg, tmp_path, *extra): + return cli.main( + ["clean", str(src), "-o", str(tmp_path / "out.csv"), "--config", str(cfg), *extra] + ) + + +def test_clean_config_privacy_detection_takes_effect(src, tmp_path, capsys): + cfg = _write_cfg( + tmp_path, {"enterprise": {"enable_privacy_detection": True, "privacy": {}}} + ) + assert _clean_loud(src, cfg, tmp_path) == 0 + assert "privacy:" in capsys.readouterr().out + assert "a@b.com" not in (tmp_path / "out.csv").read_text() + + +def test_clean_config_lineage_takes_effect(src, tmp_path): + cfg = _write_cfg(tmp_path, {"enterprise": {"lineage": {"job_name": "nightly.customers"}}}) + lineage = tmp_path / "lineage.json" + assert _clean_loud(src, cfg, tmp_path, "--quiet", "--lineage", str(lineage)) == 0 + assert "nightly.customers" in lineage.read_text() + + +def test_clean_config_k_anonymity_takes_effect(src, tmp_path, capsys): + cfg = _write_cfg( + tmp_path, + {"enterprise": {"k_anonymity": {"enabled": True, "quasi_identifiers": ["email"], "k": 2}}}, + ) + assert _clean_loud(src, cfg, tmp_path) == 0 + assert "k-anonymity (k=2)" in capsys.readouterr().out + + +def test_clean_config_entity_resolution_takes_effect(src, tmp_path, capsys): + cfg = _write_cfg( + tmp_path, + { + "enterprise": { + "enable_entity_resolution": True, + "entity_resolution": { + "backend": "pandas", + "unique_id_column": "id", + "blocking_rules": [{"sql": "l.email = r.email"}], + "comparisons": [{"column": "email"}], + }, + } + }, + ) + assert _clean_loud(src, cfg, tmp_path) == 0 + assert "entity resolution (pandas)" in capsys.readouterr().out + + +def test_clean_config_empty_clustering_object_still_means_no_clustering(): + ec = cli._build_enterprise({"clustering": {}, "enable_clustering": True}) + assert ec.clustering is None + assert ec.enable_clustering is True + + # --------------------------------------------------------------------------- # # #289: freshdata validate --suite / --contract exit 2 on unloadable rules # # --------------------------------------------------------------------------- # From 353c27333554f225096042f54b69dd5df6ddc12f Mon Sep 17 00:00:00 2001 From: Kevin Costner <120246174+kevincostner17@users.noreply.github.com> Date: Tue, 15 Sep 2026 22:16:13 +0530 Subject: [PATCH 3/5] fix(cli): accept no-op defaults for unsupported --config enterprise fields enable_contracts, drift and anonymization do nothing in freshdata clean, so the previous commit rejected them outright. That also broke configs that only spell out their defaults. A null or default value (enable_contracts: false, drift: {} or an object of DriftConfig defaults, anonymization: []) is now accepted and ignored. Any other value still exits 1 with the reason, and unknown keys inside a drift object still get a did-you-mean error. --- docs/feature-overview.md | 9 +++++--- src/freshdata/enterprise/cli.py | 25 +++++++++++++++++++--- tests/test_cli_malformed_inputs.py | 34 +++++++++++++++++++++++++++++- 3 files changed, 61 insertions(+), 7 deletions(-) diff --git a/docs/feature-overview.md b/docs/feature-overview.md index a3b9aa62..ec845461 100644 --- a/docs/feature-overview.md +++ b/docs/feature-overview.md @@ -82,10 +82,13 @@ enterprise: | `entity_resolution` | object, with `blocking_rules` and `comparisons` as lists of objects | `EntityResolutionConfig` (with `BlockingRule` and `ComparisonLevel`), applied when `enable_entity_resolution` is true | Nested objects take the field names of the class they build, and unknown names -are rejected the same way. Three `EnterpriseConfig` fields are rejected with an -explanation: `enable_contracts` and `drift` (`freshdata clean` takes no baseline or +are rejected the same way. Three `EnterpriseConfig` fields do nothing in +`freshdata clean`: `enable_contracts` and `drift` (the command takes no baseline or data contract to check against) and `anonymization` (no pipeline applies it; use -`masking`, or `privacy` with `enable_privacy_detection`). +`masking`, or `privacy` with `enable_privacy_detection`). They are accepted and +ignored when null or set to their default (`enable_contracts: false`, `drift: {}` +or an object of `DriftConfig` defaults, `anonymization: []`). Any other value is +rejected with that explanation. ## Compliance reports diff --git a/src/freshdata/enterprise/cli.py b/src/freshdata/enterprise/cli.py index 9521cfcf..93a9588a 100644 --- a/src/freshdata/enterprise/cli.py +++ b/src/freshdata/enterprise/cli.py @@ -34,6 +34,7 @@ BlockingRule, ClusterConfig, ComparisonLevel, + DriftConfig, EnterpriseConfig, EntityResolutionConfig, KAnonymityConfig, @@ -281,6 +282,21 @@ def _build_list(cls: Any, items: Any, where: str) -> tuple[Any, ...]: return tuple(_build_dataclass(cls, item, f"{where}[{i}]") for i, item in enumerate(items)) +def _is_unsupported_noop(key: str, value: Any) -> bool: + """Whether an unsupported ``enterprise`` field is null or its default, so it asks for nothing. + + Such values are accepted and ignored, so configs that spell out defaults keep loading. + """ + if value is None: + return True + if key == "enable_contracts": + return value is False + if key == "anonymization": + return value == [] + # ``drift``: an object equal to DriftConfig() configures nothing. Unknown keys still raise. + return bool(_build_dataclass(DriftConfig, value, key) == DriftConfig()) + + def _check_config_sections(data: dict[str, Any], path: str) -> None: """Reject unknown top-level sections (e.g. a misspelled ``enterprize:``).""" try: @@ -296,10 +312,13 @@ def _build_enterprise(spec: dict[str, Any]) -> EnterpriseConfig: ``EnterpriseConfig`` fields this command cannot apply, raise :class:`TypeError` instead of being silently ignored. """ + _check_keys(spec, _ENTERPRISE_KEYS | set(_ENTERPRISE_UNSUPPORTED_KEYS), "key") for key, reason in _ENTERPRISE_UNSUPPORTED_KEYS.items(): - if key in spec: - raise TypeError(f"{key!r} is not supported in --config: {reason}") - _check_keys(spec, _ENTERPRISE_KEYS, "key", suggest=_ENTERPRISE_UNSUPPORTED_KEYS) + if key in spec and not _is_unsupported_noop(key, spec[key]): + raise TypeError( + f"{key!r} is not supported in --config: {reason}; " + "only null or its default value is accepted" + ) kwargs: dict[str, Any] = {} for key in _ENTERPRISE_BOOL_KEYS: diff --git a/tests/test_cli_malformed_inputs.py b/tests/test_cli_malformed_inputs.py index 8e62a4a2..ab063c24 100644 --- a/tests/test_cli_malformed_inputs.py +++ b/tests/test_cli_malformed_inputs.py @@ -177,7 +177,12 @@ def test_clean_unknown_nested_enterprise_key_is_one_line_error( @pytest.mark.parametrize( - "key, value", [("enable_contracts", True), ("drift", {}), ("anonymization", [])] + "key, value", + [ + ("enable_contracts", True), + ("drift", {"psi_warn": 0.2}), + ("anonymization", [{"strategy": "redact"}]), + ], ) def test_clean_unsupported_enterprise_field_is_rejected(src, tmp_path, capsys, key, value): cfg = _write_cfg(tmp_path, {"enterprise": {key: value}}) @@ -185,6 +190,33 @@ def test_clean_unsupported_enterprise_field_is_rejected(src, tmp_path, capsys, k _assert_one_line_error(capsys, "cfg.json", f"'{key}' is not supported in --config") +@pytest.mark.parametrize( + "key, value", + [ + ("enable_contracts", False), + ("enable_contracts", None), + ("drift", None), + ("drift", {}), + ("drift", {"enabled": True, "psi_warn": 0.10}), + ("anonymization", []), + ("anonymization", None), + ], +) +def test_clean_unsupported_enterprise_field_at_default_is_ignored(src, tmp_path, key, value): + cfg = _write_cfg( + tmp_path, + {"enterprise": {key: value, "masking": [{"name": "m", "columns": ["email"]}]}}, + ) + assert _clean_with_config(src, cfg, tmp_path) == 0 + assert "a@b.com" not in (tmp_path / "out.csv").read_text() + + +def test_clean_default_drift_object_still_rejects_typos(src, tmp_path, capsys): + cfg = _write_cfg(tmp_path, {"enterprise": {"drift": {"psi_wrn": 0.1}}}) + assert _clean_with_config(src, cfg, tmp_path) == 1 + _assert_one_line_error(capsys, "cfg.json", "'psi_wrn' (did you mean 'psi_warn'?)") + + @pytest.mark.parametrize( "payload, needle", [ From 3982a5b535f2ce51da2f795fce1f80e3cd772da2 Mon Sep 17 00:00:00 2001 From: Kevin Costner <120246174+kevincostner17@users.noreply.github.com> Date: Tue, 15 Sep 2026 22:19:44 +0530 Subject: [PATCH 4/5] fix(cli): honour --config on native freshdata clean engines With --engine polars, duckdb, spark, freshcore or auto, cmd_clean returned into the engine path before reading --config, so the whole file was silently ignored: clean options did not apply, typos did not error, and enterprise options such as masking were dropped without a word. The config file is now loaded and validated before engine dispatch, so a typo errors on every engine. The engine path merges the 'clean' section under the command-line options, as the pandas path does, and rejects 'context'/'policy' there with a one-line error. Native engines do not run the enterprise stage, so an 'enterprise' section that sets anything other than the EnterpriseConfig defaults exits 1 and names the keys. --- docs/feature-overview.md | 7 ++ src/freshdata/enterprise/cli.py | 90 ++++++++++++++++++++----- tests/test_execution/test_cli_engine.py | 83 +++++++++++++++++++++++ 3 files changed, 163 insertions(+), 17 deletions(-) diff --git a/docs/feature-overview.md b/docs/feature-overview.md index ec845461..a2757144 100644 --- a/docs/feature-overview.md +++ b/docs/feature-overview.md @@ -90,6 +90,13 @@ ignored when null or set to their default (`enable_contracts: false`, `drift: {} or an object of `DriftConfig` defaults, `anonymization: []`). Any other value is rejected with that explanation. +The config file is validated the same way on every `--engine`. With a native +engine (`polars`, `duckdb`, `spark`, `freshcore`, `auto`) the `clean` section is +applied under the command-line options, as on pandas, except `context` and +`policy`, which only the pandas engine supports. Native engines do not run the +enterprise stage, so an `enterprise` section that sets anything other than the +defaults exits 1 and names the keys; drop them or use `--engine pandas`. + ## Compliance reports The `freshdata.compliance` subpackage turns a `CleanReport` into a regulatory diff --git a/src/freshdata/enterprise/cli.py b/src/freshdata/enterprise/cli.py index 93a9588a..0a40791f 100644 --- a/src/freshdata/enterprise/cli.py +++ b/src/freshdata/enterprise/cli.py @@ -347,6 +347,33 @@ def _build_enterprise(spec: dict[str, Any]) -> EnterpriseConfig: return EnterpriseConfig(**kwargs) +def _read_config_arg(args: argparse.Namespace) -> tuple[dict[str, Any], EnterpriseConfig]: + """Load and validate ``--config``: its ``clean`` options and its built EnterpriseConfig. + + Unknown sections or keys and invalid values raise ``ValueError`` naming the file. + Without ``--config`` this returns ``({}, EnterpriseConfig())``. + """ + if not args.config: + return {}, EnterpriseConfig() + data = _load_config_file(args.config) + _check_config_sections(data, args.config) + file_clean = _config_section(data, "clean", args.config) + try: + ec = _build_enterprise(_config_section(data, "enterprise", args.config)) + except TypeError as exc: + # Unknown/misspelled keys (MaskingRule(**rule)) or a non-object entry. + raise ValueError( + f"invalid 'enterprise' section in config file {args.config}: {exc}" + ) from exc + try: + merge_options(None, **file_clean) + except TypeError as exc: # unknown option names, e.g. a typo in the config file + raise ValueError( + f"invalid 'clean' options in config file {args.config}: {exc}" + ) from exc + return file_clean, ec + + def _load_profile_arg(path: str, *, quiet: bool = False) -> tuple[Any, int]: """Load a .fdprofile for the CLI: (profile, 0) or (None, exit_code).""" from ..learning import load_profile # noqa: PLC0415 - lazy import @@ -370,31 +397,21 @@ def _load_profile_arg(path: str, *, quiet: bool = False) -> tuple[Any, int]: def cmd_clean(args: argparse.Namespace) -> int: if getattr(args, "engine", None) and args.engine != "pandas": + # Validate --config before anything else, so a typo errors on every engine. + engine_clean, engine_ec = _read_config_arg(args) if getattr(args, "context_file", None): print("error: --context-file is only supported on the pandas engine") return 2 if getattr(args, "profile", None): print("error: --profile is only supported on the pandas engine") return 2 - return _cmd_clean_engine(args) + return _cmd_clean_engine(args, engine_clean, engine_ec) learned_profile = None if getattr(args, "profile", None): learned_profile, code = _load_profile_arg(args.profile, quiet=args.quiet) if learned_profile is None: return code - file_clean: dict[str, Any] = {} - ec = EnterpriseConfig() - if args.config: - data = _load_config_file(args.config) - _check_config_sections(data, args.config) - file_clean = _config_section(data, "clean", args.config) - try: - ec = _build_enterprise(_config_section(data, "enterprise", args.config)) - except TypeError as exc: - # Unknown/misspelled keys (MaskingRule(**rule)) or a non-object entry. - raise ValueError( - f"invalid 'enterprise' section in config file {args.config}: {exc}" - ) from exc + file_clean, ec = _read_config_arg(args) overrides: dict[str, Any] = {"strategy": args.strategy} if args.strategy else {} if getattr(args, "drop_duplicates", None): @@ -492,21 +509,60 @@ def cmd_clean(args: argparse.Namespace) -> int: return 0 if result.passed_gate else 1 -def _cmd_clean_engine(args: argparse.Namespace) -> int: +def _check_engine_enterprise(ec: EnterpriseConfig, args: argparse.Namespace) -> None: + """Raise ``ValueError`` when *ec* asks for enterprise work a native ``--engine`` skips.""" + default = EnterpriseConfig() + changed = [ + f.name + for f in dataclasses.fields(EnterpriseConfig) + if getattr(ec, f.name) != getattr(default, f.name) + ] + if changed: + raise ValueError( + f"config file {args.config} sets 'enterprise' options that --engine " + f"{args.engine} cannot run: {', '.join(changed)}. The enterprise stage " + "(masking, clustering, validation, privacy, trust gate, lineage) only runs " + "on --engine pandas; remove these keys or use --engine pandas" + ) + + +def _cmd_clean_engine( + args: argparse.Namespace, + file_clean: dict[str, Any] | None = None, + enterprise: EnterpriseConfig | None = None, +) -> int: """Clean via a scalable execution backend (polars / duckdb / spark / auto). The input path is handed straight to the backend so it can read it natively (DuckDB/Polars scan files in place; Spark reads via its own readers). The cleaned result is converted to pandas for writing and a CleanReport summary. + + A ``--config`` file's ``clean`` options are merged under the command-line + options, as on the pandas engine. The enterprise stage does not run here, so an + ``enterprise`` section that sets anything beyond the defaults is an error. """ import freshdata as fd from ..execution import EngineConfig - overrides = {"strategy": args.strategy} if args.strategy else {} + overrides: dict[str, Any] = {"strategy": args.strategy} if args.strategy else {} if getattr(args, "drop_duplicates", None): overrides["drop_duplicates"] = True - clean_config = merge_options(None, **overrides) if overrides else None + merged = {**(file_clean or {}), **overrides} + source = f" in config file {args.config}" if args.config else "" + try: + clean_config = merge_options(None, **merged) if merged else None + except TypeError as exc: + raise ValueError(f"invalid 'clean' options{source}: {exc}") from exc + if clean_config is not None and ( + clean_config.context is not None or clean_config.policy is not None + ): + raise ValueError( + f"the 'context' and 'policy' clean options{source} are only supported " + "on the pandas engine" + ) + if enterprise is not None: + _check_engine_enterprise(enterprise, args) engine_config = EngineConfig(engine=args.engine, output_format="pandas") if getattr(args, "memory_limit_gb", None) is not None: diff --git a/tests/test_execution/test_cli_engine.py b/tests/test_execution/test_cli_engine.py index 8917817c..01f6c762 100644 --- a/tests/test_execution/test_cli_engine.py +++ b/tests/test_execution/test_cli_engine.py @@ -64,3 +64,86 @@ def test_clean_cli_engine_report(parquet_in, tmp_path): with open(report_path, encoding="utf-8") as fh: payload = json.load(fh) assert "actions" in payload + + +# --------------------------------------------------------------------------- # +# --config on native engines: validated first, clean applied, enterprise gated # +# --------------------------------------------------------------------------- # +_NATIVE_ENGINES = ["polars", "duckdb", "spark", "freshcore", "auto"] + + +def _run_with_config(parquet_in, tmp_path, engine, payload, *extra): + cfg = tmp_path / "cfg.json" + cfg.write_text(json.dumps(payload)) + out_path = tmp_path / f"out_{engine}.parquet" + rc = main([ + "clean", parquet_in, "-o", str(out_path), "--engine", engine, + "--config", str(cfg), "--quiet", *extra, + ]) + return rc, out_path + + +@pytest.mark.parametrize("engine", _NATIVE_ENGINES) +@pytest.mark.parametrize( + "payload, needle", + [ + ({"clean": {"stratgy": "conservative"}}, "'stratgy' (did you mean 'strategy'?)"), + ({"enterprise": {"enable_maskin": True}}, "'enable_maskin' (did you mean"), + ({"enterprize": {}}, "'enterprize' (did you mean 'enterprise'?)"), + ], +) +def test_clean_cli_engine_config_typos_error_on_every_engine( + parquet_in, tmp_path, capsys, engine, payload, needle +): + rc, out_path = _run_with_config(parquet_in, tmp_path, engine, payload) + assert rc == 1 + captured = capsys.readouterr() + assert needle in captured.err + assert "Traceback" not in captured.err + assert not out_path.exists() + + +@pytest.mark.parametrize("engine", _NATIVE_ENGINES) +def test_clean_cli_engine_rejects_enterprise_features(parquet_in, tmp_path, capsys, engine): + payload = { + "enterprise": { + "masking": [{"name": "m", "columns": ["Customer ID"]}], + "fail_under_trust": 80, + } + } + rc, out_path = _run_with_config(parquet_in, tmp_path, engine, payload) + assert rc == 1 + err = capsys.readouterr().err + assert f"--engine {engine} cannot run: masking, fail_under_trust" in err + assert "--engine pandas" in err + assert not out_path.exists() + + +@pytest.mark.parametrize("engine", _NATIVE_ENGINES) +def test_clean_cli_engine_rejects_context_in_clean_section(parquet_in, tmp_path, capsys, engine): + rc, out_path = _run_with_config( + parquet_in, tmp_path, engine, {"clean": {"context": "never drop rows"}} + ) + assert rc == 1 + assert "only supported on the pandas engine" in capsys.readouterr().err + assert not out_path.exists() + + +@pytest.mark.parametrize("engine", ["polars", "duckdb", "spark"]) +def test_clean_cli_engine_applies_clean_section(parquet_in, tmp_path, engine): + # Before the fix the engine path never read --config, so columns were renamed. + pytest.importorskip("pyspark" if engine == "spark" else engine) + rc, out_path = _run_with_config( + parquet_in, tmp_path, engine, {"clean": {"column_names": False}} + ) + assert rc == 0 + assert "Customer ID" in pd.read_parquet(out_path).columns + + +@pytest.mark.parametrize("engine", ["polars", "duckdb"]) +def test_clean_cli_engine_accepts_default_enterprise_section(parquet_in, tmp_path, engine): + pytest.importorskip(engine) + payload = {"enterprise": {"enable_masking": True, "enable_contracts": False, "drift": None}} + rc, out_path = _run_with_config(parquet_in, tmp_path, engine, payload) + assert rc == 0 + assert "customer_id" in pd.read_parquet(out_path).columns From d1e517ab31cb65d824063ee63f8317f13d21b1e8 Mon Sep 17 00:00:00 2001 From: Kevin Costner <120246174+kevincostner17@users.noreply.github.com> Date: Tue, 15 Sep 2026 22:22:23 +0530 Subject: [PATCH 5/5] fix(cli): print exit-2 and exit-3 errors to stderr freshdata clean, profile, learn, plan, apply-plan, policy compile and models pull printed their own error and usage messages to stdout, so a pipeline reading stdout (for example --output-format json | jq) got an error line in its data. Examples: 'error: cannot load profile ...', 'error: Unknown model id ...'. Exit-1 errors and validate's exit-2 errors already went to stderr. These messages now go to stderr in the one-line 'freshdata: error: ...' form that main's top-level handler uses. Exit codes are unchanged. Tests that read these messages from stdout now read stderr. --- src/freshdata/enterprise/cli.py | 47 ++++++++++++-------- tests/context/test_context_cli.py | 4 +- tests/learning/test_cli.py | 10 ++--- tests/test_cli_error_streams.py | 73 +++++++++++++++++++++++++++++++ tests/test_cli_models.py | 15 ++++--- 5 files changed, 117 insertions(+), 32 deletions(-) create mode 100644 tests/test_cli_error_streams.py diff --git a/src/freshdata/enterprise/cli.py b/src/freshdata/enterprise/cli.py index 0a40791f..27385157 100644 --- a/src/freshdata/enterprise/cli.py +++ b/src/freshdata/enterprise/cli.py @@ -87,6 +87,15 @@ def _safe_print(text: str) -> None: print(text.encode(encoding, errors="replace").decode(encoding, errors="replace")) +def _print_error(message: str) -> None: + """Print a ``freshdata: error: ...`` diagnostic to stderr. + + Same format as the top-level handler in :func:`main`, so stdout stays reserved for + reports and JSON even when a command fails with its own exit code. + """ + print(f"freshdata: error: {message}", file=sys.stderr) + + def _emit_report(report: Any, args: argparse.Namespace, legacy_text: str) -> None: """Print a clean report honoring the display flags. @@ -382,10 +391,10 @@ def _load_profile_arg(path: str, *, quiet: bool = False) -> tuple[Any, int]: try: profile = load_profile(path) except ProfileError as exc: - print(f"error: cannot load profile {path}: {exc}") + _print_error(f"cannot load profile {path}: {exc}") return None, 2 except (OSError, ValueError) as exc: - print(f"error: cannot read profile {path}: {exc}") + _print_error(f"cannot read profile {path}: {exc}") return None, 2 if getattr(profile.manifest, "contains_raw_values", False) and not quiet: print( @@ -400,10 +409,10 @@ def cmd_clean(args: argparse.Namespace) -> int: # Validate --config before anything else, so a typo errors on every engine. engine_clean, engine_ec = _read_config_arg(args) if getattr(args, "context_file", None): - print("error: --context-file is only supported on the pandas engine") + _print_error("--context-file is only supported on the pandas engine") return 2 if getattr(args, "profile", None): - print("error: --profile is only supported on the pandas engine") + _print_error("--profile is only supported on the pandas engine") return 2 return _cmd_clean_engine(args, engine_clean, engine_ec) learned_profile = None @@ -477,7 +486,7 @@ def cmd_clean(args: argparse.Namespace) -> int: profile=learned_profile, ) except PolicyError as exc: - print(f"error: {exc}") + _print_error(str(exc)) return 2 if args.output: @@ -609,8 +618,8 @@ def cmd_profile(args: argparse.Namespace) -> int: if args.input in ("audit", "diff", "merge"): return _cmd_profile_tools(args) if getattr(args, "paths", None): - print( - f"error: unexpected extra arguments {args.paths}; " + _print_error( + f"unexpected extra arguments {args.paths}; " "did you mean 'freshdata profile audit|diff|merge'?" ) return 2 @@ -630,7 +639,7 @@ def _cmd_profile_tools(args: argparse.Namespace) -> int: paths = list(getattr(args, "paths", []) or []) if tool == "audit": if len(paths) != 1: - print("usage: freshdata profile audit PROFILE.fdprofile [--json]") + _print_error("usage: freshdata profile audit PROFILE.fdprofile [--json]") return 2 profile, code = _load_profile_arg(paths[0]) if profile is None: @@ -644,7 +653,7 @@ def _cmd_profile_tools(args: argparse.Namespace) -> int: return 1 if audit.raw_sensitive_literals else 0 if tool == "diff": if len(paths) != 2: - print("usage: freshdata profile diff A.fdprofile B.fdprofile") + _print_error("usage: freshdata profile diff A.fdprofile B.fdprofile") return 2 left, code = _load_profile_arg(paths[0]) if left is None: @@ -657,13 +666,13 @@ def _cmd_profile_tools(args: argparse.Namespace) -> int: return 0 if diff.is_empty else 1 # merge if len(paths) != 2: - print( + _print_error( "usage: freshdata profile merge A.fdprofile B.fdprofile " "-o MERGED.fdprofile [--strategy STRATEGY]" ) return 2 if not getattr(args, "output", None): - print("error: profile merge requires -o/--output for the merged profile") + _print_error("profile merge requires -o/--output for the merged profile") return 2 left, code = _load_profile_arg(paths[0]) if left is None: @@ -676,7 +685,7 @@ def _cmd_profile_tools(args: argparse.Namespace) -> int: try: merged = left.merge(right, strategy=args.strategy) except ProfileMergeError as exc: - print(f"error: {exc}") + _print_error(str(exc)) return 2 from ..learning import save_profile # noqa: PLC0415 - lazy import @@ -710,7 +719,7 @@ def cmd_learn(args: argparse.Namespace) -> int: min_precision=args.min_precision, ) except (ProfileError, ValueError, TypeError) as exc: - print(f"error: {exc}") + _print_error(str(exc)) return 2 save_profile(profile, args.output) if not args.quiet: @@ -834,7 +843,7 @@ def cmd_policy_compile(args: argparse.Namespace) -> int: try: policy = compile_context(text, columns=columns, strict=args.strict) except PolicyError as exc: - print(f"error: {exc}") + _print_error(str(exc)) return 2 print(policy.summary()) if args.output: @@ -868,7 +877,7 @@ def cmd_models_pull(args: argparse.Namespace) -> int: try: path = models.pull(args.model_id, force=args.force) except models.ModelError as exc: - print(f"error: {exc}") + _print_error(str(exc)) return 2 print(f"pulled {args.model_id} -> {path}") return 0 @@ -895,11 +904,11 @@ def cmd_plan(args: argparse.Namespace) -> int: try: plan = fd.suggest_plan(df, **_plan_overrides(args)) except PolicyError as exc: - print(f"error: {exc}") + _print_error(str(exc)) return 2 repair_plan = plan.repair_plan if repair_plan is None: - print( + _print_error( "no repair plan: pass --context-file and/or --semantic-mode so planned actions exist" ) return 2 @@ -922,10 +931,10 @@ def cmd_apply_plan(args: argparse.Namespace) -> int: try: cleaned, report = fd.apply_plan(df, plan, allow_drift=args.allow_drift) except fd.PlanDriftError as exc: - print(f"error: {exc}") + _print_error(str(exc)) return 2 except fd.ProtectedColumnError as exc: - print(f"error: {exc}") + _print_error(str(exc)) return 3 if args.output: _write_frame( diff --git a/tests/context/test_context_cli.py b/tests/context/test_context_cli.py index 2e2063e1..7442d6d3 100644 --- a/tests/context/test_context_cli.py +++ b/tests/context/test_context_cli.py @@ -63,7 +63,7 @@ def test_policy_compile_strict_fails_on_unparsed(tmp_path, data_csv, capsys): bad.write_text("Utter gibberish sentence.\n", encoding="utf-8") code = main(["policy", "compile", str(bad), "--schema", str(data_csv), "--strict"]) assert code == 2 - assert "unparsed_sentence" in capsys.readouterr().out + assert "unparsed_sentence" in capsys.readouterr().err def test_policy_compile_strict_fails_on_unresolved(tmp_path, data_csv, capsys): @@ -71,7 +71,7 @@ def test_policy_compile_strict_fails_on_unresolved(tmp_path, data_csv, capsys): bad.write_text("heart_rate must be between 0 and 200.\n", encoding="utf-8") code = main(["policy", "compile", str(bad), "--schema", str(data_csv), "--strict"]) assert code == 2 - assert "unresolved" in capsys.readouterr().out + assert "unresolved" in capsys.readouterr().err def test_clean_with_context_file(data_csv, rules_txt, tmp_path, capsys): diff --git a/tests/learning/test_cli.py b/tests/learning/test_cli.py index 336df412..6948d9e3 100644 --- a/tests/learning/test_cli.py +++ b/tests/learning/test_cli.py @@ -103,7 +103,7 @@ def test_clean_profile_corrupt_exits_nonzero(self, tmp_path, new_csv, capsys): bad.write_bytes(b"not a zip") code = cli.main(["clean", str(new_csv), "--profile", str(bad)]) assert code == 2 - assert "error" in capsys.readouterr().out + assert "error" in capsys.readouterr().err def test_clean_profile_non_pandas_engine_rejected(self, profile_path, new_csv, capsys): code = cli.main( @@ -117,7 +117,7 @@ def test_clean_profile_non_pandas_engine_rejected(self, profile_path, new_csv, c ] ) assert code == 2 - assert "pandas engine" in capsys.readouterr().out + assert "pandas engine" in capsys.readouterr().err class TestProfileTools: @@ -146,7 +146,7 @@ def test_audit_corrupt_hash_exits_nonzero(self, profile_path, tmp_path, capsys): dst.writestr(name, data) code = cli.main(["profile", "audit", str(tampered)]) assert code == 2 - assert "error" in capsys.readouterr().out.lower() + assert "error" in capsys.readouterr().err.lower() def test_diff_identical_exit_zero(self, profile_path, capsys): code = cli.main(["profile", "diff", str(profile_path), str(profile_path)]) @@ -181,7 +181,7 @@ def test_merge_writes_output(self, profile_path, tmp_path, capsys): def test_merge_requires_output(self, profile_path, capsys): code = cli.main(["profile", "merge", str(profile_path), str(profile_path)]) assert code == 2 - assert "-o" in capsys.readouterr().out + assert "-o" in capsys.readouterr().err def test_usage_errors(self, capsys): assert cli.main(["profile", "audit"]) == 2 @@ -197,4 +197,4 @@ def test_data_profiling_still_works(self, new_csv, capsys): def test_extra_args_on_data_profiling_rejected(self, new_csv, capsys): code = cli.main(["profile", str(new_csv), "extra.csv"]) assert code == 2 - assert "audit|diff|merge" in capsys.readouterr().out + assert "audit|diff|merge" in capsys.readouterr().err diff --git a/tests/test_cli_error_streams.py b/tests/test_cli_error_streams.py new file mode 100644 index 00000000..d6a7b68d --- /dev/null +++ b/tests/test_cli_error_streams.py @@ -0,0 +1,73 @@ +"""Exit-2/3 CLI errors go to stderr, never stdout (FDC-L5-010). + +Every ``freshdata`` failure uses the one-line ``freshdata: error: ...`` form on +stderr, as the top-level handler does for exit 1, so a pipeline reading stdout +(for example ``--output-format json | jq``) never gets an error line in its data. +Exit codes are unchanged. +""" + +from __future__ import annotations + +import pandas as pd +import pytest + +from freshdata.enterprise import cli + + +@pytest.fixture +def files(tmp_path): + data = tmp_path / "in.csv" + pd.DataFrame({"id": [1, 2], "v": [2, 3]}).to_csv(data, index=False) + rules = tmp_path / "rules.txt" + rules.write_text("frobnicate the wibble\n") + return { + "in": str(data), + "rules": str(rules), + "missing": str(tmp_path / "nope.fdprofile"), + "out": str(tmp_path / "out.csv"), + } + + +_CASES = { + "clean --profile missing": ( + ["clean", "{in}", "-o", "{out}", "--profile", "{missing}"], 2, "nope.fdprofile" + ), + "clean --strict unparsed context": ( + ["clean", "{in}", "-o", "{out}", "--context-file", "{rules}", "--strict"], 2, "" + ), + "clean --engine with --context-file": ( + ["clean", "{in}", "--engine", "polars", "--context-file", "{rules}"], + 2, + "--context-file is only supported on the pandas engine", + ), + "clean --engine with --profile": ( + ["clean", "{in}", "--engine", "duckdb", "--profile", "{missing}"], + 2, + "--profile is only supported on the pandas engine", + ), + "profile extra arguments": ( + ["profile", "{in}", "extra"], 2, "unexpected extra arguments" + ), + "profile audit usage": (["profile", "audit"], 2, "usage: freshdata profile audit"), + "profile audit missing": (["profile", "audit", "{missing}"], 2, "nope.fdprofile"), + "profile diff usage": (["profile", "diff", "{missing}"], 2, "usage: freshdata profile diff"), + "profile merge usage": ( + ["profile", "merge", "{missing}"], 2, "usage: freshdata profile merge" + ), + "profile merge without -o": ( + ["profile", "merge", "{missing}", "{missing}"], 2, "requires -o/--output" + ), + "policy compile --strict": (["policy", "compile", "{rules}", "--strict"], 2, ""), + "models pull unknown id": (["models", "pull", "no-such-model"], 2, "no-such-model"), +} + + +@pytest.mark.parametrize("label", list(_CASES)) +def test_error_goes_to_stderr_with_unchanged_exit_code(files, capsys, label): + argv, code, needle = _CASES[label] + assert cli.main([arg.format(**files) for arg in argv]) == code + captured = capsys.readouterr() + assert captured.out == "" + assert captured.err.startswith("freshdata: error: ") + assert needle in captured.err + assert "Traceback" not in captured.err diff --git a/tests/test_cli_models.py b/tests/test_cli_models.py index 0a782ca9..c685f97e 100644 --- a/tests/test_cli_models.py +++ b/tests/test_cli_models.py @@ -31,13 +31,16 @@ def test_models_status_without_semantic_extra(model_home, capsys): def test_models_pull_unpublished_errors_cleanly(model_home, capsys): assert main(["models", "pull", "fd-col-encoder-v1"]) == 2 - out = capsys.readouterr().out - assert "FRESHDATA_MODEL_URL_BASE" in out + captured = capsys.readouterr() + assert "FRESHDATA_MODEL_URL_BASE" in captured.err + assert captured.out == "" def test_models_pull_unknown_model(model_home, capsys): assert main(["models", "pull", "fd-nope-v9"]) == 2 - assert "Known models" in capsys.readouterr().out + captured = capsys.readouterr() + assert "Known models" in captured.err + assert captured.out == "" def test_models_pull_downloads_with_mocked_fetch(model_home, monkeypatch, capsys): @@ -66,9 +69,9 @@ def no_fetch(url, dest): # pragma: no cover - must not run monkeypatch.setattr(dl, "_fetch", no_fetch) assert main(["models", "pull", "fd-intent-v1"]) == 2 - out = capsys.readouterr().out - assert "Checksum mismatch" in out - assert "pulled" not in out + captured = capsys.readouterr() + assert "Checksum mismatch" in captured.err + assert captured.out == "" def test_clean_with_embedding_missing_model_prints_skip(model_home, tmp_path, capsys):