Skip to content

Commit 7b03498

Browse files
fix(dbt): honour on_low_score="fail" in FreshDataDbtTransform; unique audit files for same-alias models (#391)
on_low_score="fail" was ignored by FreshDataDbtTransform (#343): - run() only raised TrustGateError when fail_on_low_score=True, so with on_low_score="fail" a failing gate returned should_fail=True and the pipeline carried on. run() now raises when result.should_fail or when fail_on_low_score is set and the gate did not pass. The audit file is still written before raising. - run() takes a keyword-only raise_on_fail=True. gate_manifest calls run(raise_on_fail=False), so a failing model under on_low_score="fail" is recorded as failed and the run continues. The summary shape and the skipped/all_passed semantics are unchanged. Same-alias models overwrote each other's audit file (#344): - FreshDataDbtTransform gains audit_name, used as the audit file stem (<audit_name>_audit.json) instead of the table name. It is checked with _validate_audit_table_name when the transform is configured. - When output_dir is set, gate_manifest counts aliases across the models it gates (case-insensitively, since audit files may land on a case-insensitive filesystem). Models whose alias is shared are written to <schema>.<alias>_audit.json, or to <unique_id>_audit.json when the schema is missing or that name is still not unique. Other models keep <alias>_audit.json. An unsafe schema is rejected by the validator and recorded as that model's error. Closes #343 Closes #344
1 parent 1ac83fc commit 7b03498

3 files changed

Lines changed: 357 additions & 13 deletions

File tree

‎docs/integrations.md‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -111,10 +111,18 @@ result = FreshDataDbtTransform(
111111
model_name="analytics.orders",
112112
output_dir="target/freshdata",
113113
trust_score_threshold=80.0,
114-
fail_on_low_score=True,
114+
on_low_score="fail", # raise TrustGateError on a failing gate
115115
).run()
116116
```
117117

118+
With `on_low_score="fail"` (or the older `fail_on_low_score=True`), `run()` writes
119+
the audit file and then raises `TrustGateError` on a failing gate; pass
120+
`run(raise_on_fail=False)` to get the failing result back instead. When
121+
`dbt-gate` / `gate_manifest` gates several models that share an alias in different
122+
schemas, their audit files are named `<schema>.<alias>_audit.json` (or
123+
`<unique_id>_audit.json` when the schema is missing), so no model's audit
124+
overwrites another's. Other models keep `<alias>_audit.json`.
125+
118126
A bundled Jinja macro, `freshdata_trust_gate`, documents the recommended `on-run-end`
119127
invocation; see `freshdata/integrations/dbt/macros/freshdata_trust_gate.sql`.
120128

‎src/freshdata/integrations/dbt/__init__.py‎

Lines changed: 56 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
import json
2020
import logging
2121
import os
22+
from collections import Counter
2223
from dataclasses import dataclass
2324
from pathlib import Path, PurePath
2425
from typing import TYPE_CHECKING, Any
@@ -99,10 +100,15 @@ class FreshDataDbtTransform:
99100
clean_config: CleanConfig | None = None
100101
system_actor: str = "freshdata"
101102
fail_on_low_score: bool = False
103+
#: File stem for the audit (``<audit_name>_audit.json``); defaults to the table
104+
#: name. :func:`gate_manifest` sets it when two gated models share an alias.
105+
audit_name: str | None = None
102106

103107
def __post_init__(self) -> None:
104-
"""Reject invalid gate policies when the transform is configured."""
108+
"""Reject invalid gate policies and unsafe audit names at configuration."""
105109
self.on_low_score = validate_on_low_score(self.on_low_score)
110+
if self.audit_name is not None:
111+
self.audit_name = _validate_audit_table_name(self.audit_name)
106112

107113
def _split_table(self) -> tuple[str | None, str]:
108114
if self.schema:
@@ -112,24 +118,32 @@ def _split_table(self) -> tuple[str | None, str]:
112118
return ".".join(prefix) or None, table
113119
return None, self.model_name
114120

121+
def _audit_stem(self, table: str) -> str:
122+
return self.audit_name if self.audit_name is not None else table
123+
115124
def _write_audit(self, table: str, result: TrustGateResult) -> Path:
116-
table = _validate_audit_table_name(table)
125+
stem = _validate_audit_table_name(self._audit_stem(table))
117126
out_dir = Path(self.output_dir) # type: ignore[arg-type]
118127
out_dir.mkdir(parents=True, exist_ok=True)
119-
path = out_dir / f"{table}_audit.json"
128+
path = out_dir / f"{stem}_audit.json"
120129
path.write_text(json.dumps(result.to_dict(), indent=2, default=str))
121130
return path
122131

123-
def run(self) -> TrustGateResult:
124-
"""Read the model table, gate it, optionally write an audit, return the result."""
132+
def run(self, *, raise_on_fail: bool = True) -> TrustGateResult:
133+
"""Read the model table, gate it, optionally write an audit, return the result.
134+
135+
A failing gate raises :class:`TrustGateError` (after the audit is written) when
136+
``on_low_score="fail"`` or ``fail_on_low_score=True``. Pass
137+
``raise_on_fail=False`` to get the failing result back instead.
138+
"""
125139
conn = self.conn_str or os.environ.get("FRESHDATA_WAREHOUSE_CONN")
126140
if not conn:
127141
raise ValueError(
128142
"No warehouse connection: pass conn_str or set FRESHDATA_WAREHOUSE_CONN."
129143
)
130144
schema, table = self._split_table()
131145
if self.output_dir:
132-
_validate_audit_table_name(table)
146+
_validate_audit_table_name(self._audit_stem(table))
133147
df = _read_table(conn, schema, table)
134148
_, result = evaluate_trust_gate(
135149
df,
@@ -141,11 +155,39 @@ def run(self) -> TrustGateResult:
141155
)
142156
if self.output_dir:
143157
self._write_audit(table, result)
144-
if self.fail_on_low_score and not result.passed:
158+
if raise_on_fail and (
159+
result.should_fail or (self.fail_on_low_score and not result.passed)
160+
):
145161
raise TrustGateError(result.message)
146162
return result
147163

148164

165+
def _audit_names(models: list[tuple[str, dict[str, Any]]]) -> list[str | None]:
166+
"""Return an audit file stem per model, or ``None`` to keep the table name.
167+
168+
Audit files are named after the model's alias, which dbt only requires to be
169+
unique within a schema. When gated models share an alias (compared
170+
case-insensitively, as audit files may land on a case-insensitive filesystem),
171+
each of them is named ``"<schema>.<alias>"`` instead, or after its manifest
172+
``unique_id`` when it has no schema or that name is still not unique.
173+
"""
174+
tables = [node.get("alias") or node.get("name") for _, node in models]
175+
alias_counts = Counter(t.casefold() for t in tables if isinstance(t, str))
176+
names: list[str | None] = []
177+
for (node_id, node), table in zip(models, tables):
178+
if not isinstance(table, str) or alias_counts[table.casefold()] < 2:
179+
names.append(None)
180+
continue
181+
schema = node.get("schema")
182+
names.append(f"{schema}.{table}" if isinstance(schema, str) and schema else node_id)
183+
stems = [name if name is not None else table for name, table in zip(names, tables)]
184+
stem_counts = Counter(s.casefold() for s in stems if isinstance(s, str))
185+
return [
186+
node_id if isinstance(stem, str) and stem_counts[stem.casefold()] > 1 else name
187+
for (node_id, _), name, stem in zip(models, names, stems)
188+
]
189+
190+
149191
def gate_manifest(
150192
manifest_path: str | Path,
151193
*,
@@ -178,9 +220,9 @@ def gate_manifest(
178220
if not isinstance(nodes, dict):
179221
raise ValueError(f"{manifest_path} is not a dbt manifest: no 'nodes' mapping")
180222

181-
models: list[Any] = [] # raw manifest nodes (untyped JSON)
223+
models: list[tuple[str, Any]] = [] # (unique_id, raw manifest node)
182224
skipped: list[dict[str, Any]] = []
183-
for node in nodes.values():
225+
for node_id, node in nodes.items():
184226
if not isinstance(node, dict) or node.get("resource_type") != "model":
185227
continue
186228
config = node.get("config")
@@ -190,11 +232,12 @@ def gate_manifest(
190232
elif config.get("enabled") is False:
191233
skipped.append({"model": node.get("name"), "reason": "disabled"})
192234
else:
193-
models.append(node)
235+
models.append((node_id, node))
194236

237+
audit_names = _audit_names(models) if output_dir else [None] * len(models)
195238
summaries: list[dict[str, Any]] = []
196239
failed = 0
197-
for node in models:
240+
for (_, node), audit_name in zip(models, audit_names):
198241
name = node.get("name")
199242
schema = node.get("schema")
200243
table = node.get("alias") or name
@@ -208,7 +251,8 @@ def gate_manifest(
208251
output_dir=output_dir,
209252
clean_config=clean_config,
210253
system_actor=system_actor,
211-
).run()
254+
audit_name=audit_name,
255+
).run(raise_on_fail=False)
212256
except Exception as exc: # noqa: BLE001 - one bad model must not abort the run
213257
logger.warning("freshdata: gating model %r failed: %s", name, exc)
214258
summaries.append({"model": name, "error": str(exc)})

0 commit comments

Comments
 (0)