Skip to content
6 changes: 1 addition & 5 deletions tests/test_cli_inspect_paths.py
Original file line number Diff line number Diff line change
Expand Up @@ -53,10 +53,6 @@ def test_inspect_paths_reports_dataset_repo_layout_from_other_cwd(
f"clean_output: {project_example / '_smoke_out' / 'data' / 'clean' / 'project_example' / '2022' / 'project_example_2022_clean.parquet'}"
in result.output
)
assert (
f"clean_validation: {project_example / '_smoke_out' / 'data' / 'clean' / 'project_example' / '2022' / '_validate' / 'clean_validation.json'}"
in result.output
)
assert "raw_hints:" in result.output
assert "primary_output_file:" in result.output
assert "suggested_read_exists: False" in result.output
Expand Down Expand Up @@ -87,7 +83,7 @@ def test_inspect_paths_json_is_notebook_friendly(
assert payload["year"] == 2022
assert payload["config_path"] == str(config_path)
assert payload["paths"]["clean"]["output"].endswith("project_example_2022_clean.parquet")
assert payload["paths"]["clean"]["validation"].endswith("clean_validation.json")
assert payload["paths"]["clean"]["validation"] is None
assert payload["paths"]["raw"]["metadata"].endswith("metadata.json")
assert payload["paths"]["mart"]["outputs"]
assert payload["paths"]["mart"]["metadata"].endswith("metadata.json")
Expand Down
45 changes: 4 additions & 41 deletions tests/test_mcp_toolkit_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -183,66 +183,29 @@ def test_review_readiness_enriched_layers_shape(
raw_dir = root / "data" / "raw" / dataset / str(year)
raw_dir.mkdir(parents=True, exist_ok=True)
(raw_dir / "data.csv").write_bytes(b"a;b\n1;2\n")
(raw_dir / "raw_validation.json").write_text(
'{"ok":true,"errors":[],"warnings":["test warning raw"],"summary":{}}', encoding="utf-8"
)

clean_dir = root / "data" / "clean" / dataset / str(year)
clean_dir.mkdir(parents=True, exist_ok=True)
_write_parquet(clean_dir / f"{dataset}_{year}_clean.parquet")
clean_val_dir = clean_dir / "_validate"
clean_val_dir.mkdir(parents=True, exist_ok=True)
(clean_val_dir / "clean_validation.json").write_text(
json.dumps(
{
"ok": True,
"errors": [],
"warnings": ["[transition:clean] columns removed: [col_a]"],
"summary": {
"stats": {"clean_rows": 1, "clean_cols": 1, "raw_rows": 2, "row_drop_pct": 50.0}
},
"sections": {"transition": {"raw_row_count": 2, "clean_row_count": 1}},
}
),
encoding="utf-8",
)

mart_dir = root / "data" / "mart" / dataset / str(year)
mart_dir.mkdir(parents=True, exist_ok=True)
_write_parquet(mart_dir / "mart_t.parquet")
mart_val_dir = mart_dir / "_validate"
mart_val_dir.mkdir(parents=True, exist_ok=True)
(mart_val_dir / "mart_validation.json").write_text(
json.dumps(
{
"ok": False,
"errors": ["[mart_t] row_count too small"],
"warnings": [],
"summary": {"row_counts": {"mart_t": 1}},
}
),
encoding="utf-8",
)

payload = review_readiness(str(config_path), year)
layers = payload.get("layers", {})
assert "raw" in layers and "clean" in layers and "mart" in layers

# validation_msgs viene dal run record (nessun run in questo test → vuoto)
raw_msgs = layers["raw"].get("validation_msgs", {})
assert "test warning raw" in raw_msgs.get("warnings", [])
assert isinstance(raw_msgs.get("warnings", []), list)

clean_msgs = layers["clean"].get("validation_msgs", {})
assert len(clean_msgs.get("warnings", [])) == 1
assert "columns removed" in clean_msgs["warnings"][0]

trans = layers["clean"].get("transition", {})
assert trans.get("row_drop_pct") == 50.0
assert isinstance(clean_msgs.get("warnings", []), list)

mart_msgs = layers["mart"].get("validation_msgs", {})
assert len(mart_msgs.get("errors", [])) == 1
assert "row_count too small" in mart_msgs["errors"][0]
assert isinstance(mart_msgs.get("errors", []), list)

assert layers["mart"].get("validation", {}).get("ok") is False
assert isinstance(layers["raw"].get("profile", {}), dict)


Expand Down
9 changes: 3 additions & 6 deletions tests/test_validate_layers.py
Original file line number Diff line number Diff line change
Expand Up @@ -178,7 +178,7 @@ def test_check_transitions_warns_on_row_drop_over_threshold_and_removed_columns(
assert len(report["warnings"]) == 2
assert report["profiles_count"] == 1
assert any("row drop 30.0%" in warning for warning in report["warning_messages"])
assert any("columns removed from clean" in warning for warning in report["warning_messages"])
assert any("net drop" in warning for warning in report["warning_messages"])
assert any(item["kind"] == "row_drop_pct" for item in report["warnings"])
assert any(item["kind"] == "removed_columns" for item in report["warnings"])

Expand Down Expand Up @@ -263,7 +263,7 @@ def test_run_mart_validation_merges_transition_warnings_into_report(tmp_path: Pa
warnings = summary.get("warnings") or []
assert len(warnings) == 2
assert any("row drop 30.0%" in warning for warning in warnings)
assert any("columns removed from clean" in warning for warning in warnings)
assert any("net drop" in warning for warning in warnings)
sections = summary.get("sections") or {}
transition = sections.get("transition") or {}
assert transition.get("profiles_count") == 1
Expand Down Expand Up @@ -393,10 +393,7 @@ def test_run_clean_validation_raw_probe_source_legacy_autodetect(tmp_path: Path)

# With no profile, the probe must fall back to legacy autodetect
assert result["stats"].get("raw_probe_source") == "legacy_autodetect"
# Warning must mention the fallback reason (from return value, not disk)
warnings = result.get("warnings") or []
warning_text = " ".join(warnings)
assert "falling back to read_csv(auto_detect=true)" in warning_text
# Fallback message now goes to logger.debug, not to warnings


@pytest.mark.policy
Expand Down
79 changes: 9 additions & 70 deletions toolkit/clean/validate.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,6 @@
check_transitions,
required_columns_check,
)
from toolkit.quality.pa_csv_quality import assess_quality


def _clean_validation_spec(
Expand Down Expand Up @@ -406,9 +405,10 @@ def _to_snake(n: str) -> str:
clean_cols_set = set(clean_cols)
unmapped = sorted(scaffold_cols - clean_cols_set)
if unmapped:
merged_warnings.append(
f"[scaffold] {len(unmapped)} colonne raw non mappate nel clean "
f"(drop senza -- DROP: <motivo>?): {unmapped}"
logger.info(
"[scaffold] %d colonne raw non mappate nel clean: %s",
len(unmapped),
unmapped,
)
else:
profile_parse_error = True
Expand Down Expand Up @@ -445,9 +445,9 @@ def _to_snake(n: str) -> str:
_query = f"DESCRIBE SELECT * FROM read_csv('{sql_path(_raw_file)}', auto_detect=true)"
raw_probe_source = "legacy_autodetect"
if raw_probe_reason:
merged_warnings.append(
f"[scaffold] falling back to read_csv(auto_detect=true) — "
f"reason: {raw_probe_reason}. Run 'toolkit run raw -c <config>' to generate a profile."
logger.info(
"[scaffold] falling back to read_csv(auto_detect=true) — %s",
raw_probe_reason,
)
_col_rows = _con.execute(_query).fetchall()
_actual_raw_col_names = [str(r[0]) for r in _col_rows]
Expand Down Expand Up @@ -475,67 +475,10 @@ def _to_snake(n: str) -> str:

if actual_raw_col_count is None:
raw_probe_source = "unavailable"
merged_warnings.append(
"[scaffold] Profilo raw non disponibile — impossibile verificare coverage colonne raw. "
"Considera eseguire 'toolkit run raw -c <config>' per generare il profilo."
logger.info(
"[scaffold] Profilo raw non disponibile — impossibile verificare coverage colonne raw."
)

# ── PAQA quality score: valuta la qualità del CSV raw ──────────────────
# Il risultato viene aggiunto alle stats del run record per monitoraggio
# nel tempo. Usa un campione (primi 15MB) per performance.
paqa_score: int | None = None
paqa_verdict: str | None = None
paqa_semantic: int | None = None
paqa_sampled: bool = False
try:
# Legge il file CSV effettivamente usato da clean (da clean metadata)
# invece del primo *.csv alfabetico — evita di processare un versioned
# backup (file_1.csv, file_2.csv) al posto del file originale.
_clean_meta_path = out_dir / METADATA
_csv_path: Path | None = None
if _clean_meta_path.exists():
_clean_meta = read_json_or_none(_clean_meta_path)
if _clean_meta:
_inputs = (_clean_meta.get("outputs") or []) + (
_clean_meta.get("input_files") or []
)
for _f in _inputs:
_p = Path(_f) if isinstance(_f, str) else None
if _p and _p.suffix == ".csv" and _p.exists():
_csv_path = _p
break
if _csv_path is None:
_csv_files = sorted(raw_dir.glob("*.csv"))
if _csv_files:
# Fallback: primo CSV alfabetico (meno preciso)
_csv_path = _csv_files[0]
if _csv_path:
_size = _csv_path.stat().st_size
# Leggi campione: primi 15MB (stessa soglia CI sample-bytes)
_sample_bytes = min(_size, 15_728_640)
_csv_text = _csv_path.read_bytes()[:_sample_bytes].decode(
raw_profile.get("encoding_suggested", "utf-8") if raw_profile else "utf-8",
errors="replace",
)
_paqa_result = assess_quality(
_csv_text,
sampled=(_size > _sample_bytes),
known_sep=raw_profile.get("delim_suggested") if raw_profile else None,
known_encoding=raw_profile.get("encoding_suggested") if raw_profile else None,
known_skip=raw_profile.get("skip_suggested") if raw_profile else None,
)
paqa_score = _paqa_result.structural_score
paqa_verdict = _paqa_result.verdict
paqa_semantic = _paqa_result.semantic_score
paqa_sampled = _paqa_result.sampled
if _paqa_result.critical_fail:
merged_warnings.append(
f"[paqa] Qualita' CSV critica ({paqa_verdict}, score={paqa_score}): "
f"{'; '.join(_paqa_result.flags[:5])}"
)
except Exception as _paqa_err:
merged_warnings.append(f"[paqa] Quality assessment skipped: {_paqa_err}")

row_drop_pct = (
round((raw_row_count - clean_row_count) / raw_row_count * 100, 2)
if raw_row_count and clean_row_count is not None and raw_row_count > 0
Expand Down Expand Up @@ -573,10 +516,6 @@ def _to_snake(n: str) -> str:
),
**({"raw_missing_columns": raw_missing_columns} if raw_missing_columns else {}),
**({"raw_probe_source": raw_probe_source} if raw_probe_source else {}),
**({"paqa_score": paqa_score} if paqa_score is not None else {}),
**({"paqa_verdict": paqa_verdict} if paqa_verdict is not None else {}),
**({"paqa_semantic": paqa_semantic} if paqa_semantic is not None else {}),
**({"paqa_sampled": paqa_sampled} if paqa_sampled else {}),
},
"columns": clean_cols,
**({"rules": rules} if rules else {}),
Expand Down
46 changes: 39 additions & 7 deletions toolkit/contracts/pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -109,7 +109,7 @@
"extractors": _EXTRACTOR_TYPES,
"validation": {
"profile": {
"description": "Il profilo raw (raw_profile.json) rileva automaticamente encoding, delim, decimal, skip e colonne del CSV.",
"description": "Il profilo raw (raw_profile.json) rileva automaticamente encoding, delim, decimal, skip, colonne e row_count del CSV.",
"known_issue": "La profilazione potrebbe suggerire decimal='.' anche se il CSV usa ','. Va sovrascritto in clean.read.decimal.",
},
},
Expand Down Expand Up @@ -160,10 +160,22 @@
},
"transition": {
"description": (
"Il monitor di transizione confronta colonne raw vs clean. "
"Avvisa se colonne raw spariscono dal clean senza commento "
"-- DROP: <motivo>."
"Il monitor di transizione confronta raw vs clean. "
"Scatta solo se c'e' un **drop netto** di colonne "
"(rimosse - aggiunte > 0). Rinomine e selezioni "
"non generano falsi positivi."
),
"configurable_via_dataset_yml": {
"clean.validate.promotion.max_row_drop_pct": (
"Soglia % di righe perse (default: None = disabilitato). "
"Es: 15.0 = warning se si perde >15% righe raw→clean."
),
"clean.validate.promotion.warn_removed_columns": (
"Attiva/disattiva warning colonne rimosse. "
"Default: true. "
"Il warning scatta solo se net drop > 0."
),
},
},
},
"read_params": {
Expand Down Expand Up @@ -201,6 +213,22 @@
"required_tables": {
"scope": "mart.required_tables verifica che le tabelle dichiarate siano state prodotte.",
},
"transition": {
"description": (
"Monitor di transizione clean→mart. "
"Disabilitato per default (clean→mart seleziona colonne "
"di proposito). Attivabile con mart.validate.transition."
),
"configurable_via_dataset_yml": {
"mart.validate.transition.max_row_drop_pct": (
"Soglia % di righe perse tra clean e mart. Default: None = disabilitato."
),
"mart.validate.transition.warn_removed_columns": (
"Default: false. Imposta a true per vedere le "
"colonne clean non incluse nel mart."
),
},
},
},
"example_file": "project-example/sql/mart/mart_regione_anno.sql",
}
Expand All @@ -212,20 +240,20 @@
"name": "RAW",
"description": "Download file originale dalla fonte. Profilo: encoding, delim, decimal, colonne.",
"output": "CSV/parquet in data/raw/<dataset>/<anno>/",
"validation": "raw_validation.json",
"validation": "inline nel run record (_runs/)",
},
{
"name": "CLEAN",
"description": "Trasformazione SQL (clean.sql) su raw_input. Output parquet normalizzato.",
"output": "Parquet in data/clean/<dataset>/<anno>/",
"validation": "_validate/clean_validation.json",
"validation": "inline nel run record (_runs/)",
"view": RAW_INPUT_VIEW,
},
{
"name": "MART",
"description": "Aggregazione SQL (mart.sql) su clean_input. Output parquet per data-explorer/notebook.",
"output": "Parquet in data/mart/<dataset>/<anno>/",
"validation": "_validate/mart_validation.json",
"validation": "inline nel run record (_runs/)",
"view": CLEAN_INPUT_VIEW,
},
],
Expand Down Expand Up @@ -297,9 +325,12 @@
"clean.read.encoding (encoding CSV, default utf-8, PA spesso cp1252)",
"clean.read.decimal (separatore decimale, default '.'; usa ',' per italiano)",
"clean.required_columns (lista colonne OUTPUT attese nel clean)",
"clean.validate.promotion.warn_removed_columns (attiva warning colonne rimosse raw→clean, default true)",
"clean.validate.promotion.max_row_drop_pct (soglia % righe perse raw→clean, default none)",
"mart.tables (lista tabelle MART con nome e path SQL)",
"mart.tables[].years (per tabelle multi-anno)",
"mart.validate.table_rules (regole per tabella: primary_key, not_null, ranges)",
"mart.validate.transition.max_row_drop_pct (soglia % righe perse clean→mart, default none)",
"raw.sources[].type (http_file, ckan, sdmx, sparql, local_file)",
"raw.sources[].extractor (identity, unzip_all, unzip_first, unzip_first_csv)",
"support (lista dataset di supporto, eseguiti prima del candidate)",
Expand Down Expand Up @@ -348,6 +379,7 @@
"required_columns = nomi OUTPUT del clean, non raw | "
"se decimal=',' basta CAST(x AS DOUBLE) | "
"mart.sql: SELECT ... FROM clean_input | "
"validazione: inline nel run record (_runs/), non piu' file separati | "
"comandi: toolkit run init / preflight / all / scout / inspect"
),
}
8 changes: 6 additions & 2 deletions toolkit/core/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -351,7 +351,9 @@ def from_dict(d: dict) -> MartTableConfig:
@dataclass
class MartValidateConfig:
table_rules: dict[str, MartTableRuleConfig] = field(default_factory=dict)
transition: TransitionConfig = field(default_factory=TransitionConfig)
transition: TransitionConfig = field(
default_factory=lambda: TransitionConfig(warn_removed_columns=False)
)

@staticmethod
def from_dict(d: dict | None) -> MartValidateConfig | None:
Expand All @@ -365,7 +367,9 @@ def from_dict(d: dict | None) -> MartValidateConfig | None:
rules[k] = v
trans = d.get("transition") or d.get("transition")
trans_obj = (
TransitionConfig(**trans) if trans and isinstance(trans, dict) else TransitionConfig()
TransitionConfig(**trans)
if trans and isinstance(trans, dict)
else TransitionConfig(warn_removed_columns=False)
)
return MartValidateConfig(table_rules=rules, transition=trans_obj)

Expand Down
5 changes: 2 additions & 3 deletions toolkit/core/paths.py
Original file line number Diff line number Diff line change
Expand Up @@ -75,12 +75,11 @@ def serialize_metadata_path(path: Path | None, rel_root: Path | None) -> str | N
# ---------------------------------------------------------------------------

# Validation
RAW_VALIDATION = "raw_validation.json"
CLEAN_VALIDATION = "_validate/clean_validation.json"
MART_VALIDATION = "_validate/mart_validation.json"

# Profile (raw only)
RAW_PROFILE_DIR = "_profile"


RAW_PROFILE = "raw_profile.json" # sotto _profile/
RAW_SUGGESTED_READ = "suggested_read.yml" # sotto _profile/

Expand Down
Loading
Loading